diff --git a/backend/src/main.rs b/backend/src/main.rs index 6c851461a8..c1141246cd 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -23,9 +23,10 @@ use tokio::{ use windmill_api::HTTP_CLIENT; use windmill_common::{ global_settings::{ - BASE_URL_SETTING, CUSTOM_TAGS_SETTING, DISABLE_STATS_SETTING, ENV_SETTINGS, EXPOSE_METRICS, - EXTRA_PIP_INDEX_URL_SETTING, LICENSE_KEY_SETTING, NPM_CONFIG_REGISTRY_SETTING, - OAUTH_SETTING, REQUEST_SIZE_LIMIT_SETTING, RETENTION_PERIOD_SECS_SETTING, + BASE_URL_SETTING, CUSTOM_TAGS_SETTING, DISABLE_STATS_SETTING, ENV_SETTINGS, + EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING, + KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, + REQUEST_SIZE_LIMIT_SETTING, RETENTION_PERIOD_SECS_SETTING, }, stats::schedule_stats, utils::rd_string, @@ -40,9 +41,9 @@ use windmill_worker::{ }; use crate::monitor::{ - initial_load, monitor_db, reload_base_url_setting, reload_extra_pip_index_url_setting, - reload_license_key, reload_npm_config_registry_setting, reload_retention_period_setting, - reload_server_config, reload_worker_config, + initial_load, load_keep_job_dir, monitor_db, reload_base_url_setting, + reload_extra_pip_index_url_setting, reload_license_key, reload_npm_config_registry_setting, + reload_retention_period_setting, reload_server_config, reload_worker_config, }; const GIT_VERSION: &str = git_version!(args = ["--tag", "--always"], fallback = "unknown-version"); @@ -341,7 +342,10 @@ Windmill Community Edition {GIT_VERSION} NPM_CONFIG_REGISTRY_SETTING => { reload_npm_config_registry_setting(&db).await }, - EXPOSE_METRICS => { + KEEP_JOB_DIR_SETTING => { + load_keep_job_dir(&db).await; + } + EXPOSE_METRICS_SETTING | EXPOSE_DEBUG_METRICS_SETTING => { tracing::info!("Metrics setting changed, restarting"); // we wait a bit randomly to avoid having all serverss and workers shutdown at same time let rd_delay = rand::thread_rng().gen_range(0..4); diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 72a623f928..08aa850e6d 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -21,7 +21,8 @@ use windmill_api::{ use windmill_common::{ error, global_settings::{ - BASE_URL_SETTING, EXPOSE_METRICS, EXTRA_PIP_INDEX_URL_SETTING, LICENSE_KEY_SETTING, + BASE_URL_SETTING, EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, + EXTRA_PIP_INDEX_URL_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, REQUEST_SIZE_LIMIT_SETTING, RETENTION_PERIOD_SECS_SETTING, }, @@ -29,10 +30,10 @@ use windmill_common::{ server::load_server_config, users::truncate_token, worker::{load_worker_config, reload_custom_tags_setting, SERVER_CONFIG, WORKER_CONFIG}, - BASE_URL, DB, METRICS_ENABLED, + BASE_URL, DB, METRICS_DEBUG_ENABLED, METRICS_ENABLED, }; use windmill_worker::{ - create_token_for_owner, handle_job_error, AuthedClient, NPM_CONFIG_REGISTRY, + create_token_for_owner, handle_job_error, AuthedClient, KEEP_JOB_DIR, NPM_CONFIG_REGISTRY, PIP_EXTRA_INDEX_URL, SCRIPT_TOKEN_EXPIRY, }; @@ -84,80 +85,60 @@ pub async fn initial_load( server_mode: bool, ) { if let Err(e) = load_metrics_enabled(db).await { - tracing::error!("Error reloading loading metrics: {e}"); + tracing::error!("Error loading expose metrics: {e}"); } - let reload_worker_config_f = async { - if worker_mode { - reload_worker_config(&db, tx, false).await; - } - }; - let reload_custom_tags_f = async { - if server_mode { - if let Err(e) = reload_custom_tags_setting(db).await { - tracing::error!("Error reloading custom tags: {:?}", e) - } - } - }; - let reload_base_url_f = async { - if let Err(e) = reload_base_url_setting(db).await { - tracing::error!("Error reloading base url: {:?}", e) - } - }; + if let Err(e) = load_metrics_debug_enabled(db).await { + tracing::error!("Error loading expose debug metrics: {e}"); + } - let reload_server_config_f = async { - if server_mode { - reload_server_config(&db).await; - } - }; - let reload_retention_period_f = async { - if server_mode { - reload_retention_period_setting(&db).await; - } - }; + if worker_mode { + load_keep_job_dir(db).await; + } - let reload_request_size_f = async { - if server_mode { - reload_request_size(&db).await; - } - }; + if worker_mode { + reload_worker_config(&db, tx, false).await; + } - let reload_license_key_f = async { - #[cfg(feature = "enterprise")] - if let Err(e) = reload_license_key(&db).await { - tracing::error!("Error reloading license key: {:?}", e) + if server_mode { + if let Err(e) = reload_custom_tags_setting(db).await { + tracing::error!("Error reloading custom tags: {:?}", e) } - }; + } - let reload_extra_pip_index_url_f = async { - if worker_mode { - reload_extra_pip_index_url_setting(&db).await; - } - }; + if let Err(e) = reload_base_url_setting(db).await { + tracing::error!("Error reloading base url: {:?}", e) + } - let reload_npm_config_registry_f = async { - if worker_mode { - reload_npm_config_registry_setting(&db).await; - } - }; + if server_mode { + reload_server_config(&db).await; + } - join!( - reload_worker_config_f, - reload_server_config_f, - reload_custom_tags_f, - reload_request_size_f, - reload_base_url_f, - reload_retention_period_f, - reload_license_key_f, - reload_extra_pip_index_url_f, - reload_npm_config_registry_f - ); + if server_mode { + reload_retention_period_setting(&db).await; + } + if server_mode { + reload_request_size(&db).await; + } + + #[cfg(feature = "enterprise")] + if let Err(e) = reload_license_key(&db).await { + tracing::error!("Error reloading license key: {:?}", e) + } + + if worker_mode { + reload_extra_pip_index_url_setting(&db).await; + } + + if worker_mode { + reload_npm_config_registry_setting(&db).await; + } } pub async fn load_metrics_enabled(db: &DB) -> error::Result<()> { let metrics_enabled = sqlx::query_scalar!( "SELECT value FROM global_settings WHERE name = $1", - EXPOSE_METRICS + EXPOSE_METRICS_SETTING ) .fetch_optional(db) .await; @@ -168,6 +149,35 @@ pub async fn load_metrics_enabled(db: &DB) -> error::Result<()> { Ok(()) } +pub async fn load_metrics_debug_enabled(db: &DB) -> error::Result<()> { + let metrics_enabled = sqlx::query_scalar!( + "SELECT value FROM global_settings WHERE name = $1", + EXPOSE_DEBUG_METRICS_SETTING + ) + .fetch_optional(db) + .await; + match metrics_enabled { + Ok(Some(serde_json::Value::Bool(t))) => METRICS_DEBUG_ENABLED.store(t, Ordering::Relaxed), + _ => (), + }; + Ok(()) +} +pub async fn load_keep_job_dir(db: &DB) { + let metrics_enabled = sqlx::query_scalar!( + "SELECT value FROM global_settings WHERE name = $1", + KEEP_JOB_DIR_SETTING + ) + .fetch_optional(db) + .await; + match metrics_enabled { + Ok(Some(serde_json::Value::Bool(t))) => KEEP_JOB_DIR.store(t, Ordering::Relaxed), + Err(e) => { + tracing::error!("Error loading keep job dir metrics: {e}"); + } + _ => (), + }; +} + pub async fn delete_expired_items(db: &DB) -> () { let tokens_deleted_r: std::result::Result, _> = sqlx::query_scalar( "DELETE FROM token WHERE expiration <= now() diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 3ed0b0ba07..06c15fa8a0 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -9,7 +9,9 @@ pub const NPM_CONFIG_REGISTRY_SETTING: &str = "npm_config_registry"; pub const EXTRA_PIP_INDEX_URL_SETTING: &str = "pip_extra_index_url"; pub const UNIQUE_ID_SETTING: &str = "uid"; pub const DISABLE_STATS_SETTING: &str = "disable_stats"; -pub const EXPOSE_METRICS: &str = "expose_metrics"; +pub const EXPOSE_METRICS_SETTING: &str = "expose_metrics"; +pub const EXPOSE_DEBUG_METRICS_SETTING: &str = "expose_debug_metrics"; +pub const KEEP_JOB_DIR_SETTING: &str = "keep_job_dir"; pub const ENV_SETTINGS: [&str; 54] = [ "DISABLE_NSJAIL", diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 45d5e97fc7..1c1097c1cb 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -59,6 +59,8 @@ lazy_static::lazy_static! { .unwrap_or_else(|| SocketAddr::from(([0, 0, 0, 0], *METRICS_PORT))); pub static ref METRICS_ENABLED: AtomicBool = AtomicBool::new(std::env::var("METRICS_PORT").is_ok() || std::env::var("METRICS_ADDR").is_ok()); + pub static ref METRICS_DEBUG_ENABLED: AtomicBool = AtomicBool::new(false); + pub static ref BASE_URL: Arc> = Arc::new(RwLock::new("".to_string())); pub static ref IS_READY: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false); } diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 1cae723c27..32193850c5 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -127,7 +127,7 @@ pub async fn cancel_job<'c: 'async_recursion>( } else { let reason = reason .clone() - .unwrap_or_else(|| "No reason provided".to_string()); + .unwrap_or_else(|| "unexplicited reasons".to_string()); let e = serde_json::json!({"message": format!("Job canceled: {reason} by {username}"), "name": "Canceled", "reason": reason, "canceler": username}); let add_job = add_completed_job_error( &db, diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 741ad22ae7..1c7d1e555c 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -483,7 +483,7 @@ pub async fn handle_child( if current_mem > *mem_peak { *mem_peak = current_mem } - tracing::info!("{job_id} in {_w_id} still running. mem: {current_mem}kB, peak mem: {mem_peak}kB"); + tracing::info!("{worker_name}/{job_id} in {_w_id} still running. mem: {current_mem}kB, peak mem: {mem_peak}kB"); if sqlx::query_scalar!("UPDATE queue SET mem_peak = $1, last_ping = now() WHERE id = $2 RETURNING canceled", *mem_peak, job_id) .fetch_optional(&db) .await @@ -693,7 +693,7 @@ pub async fn handle_child( let (wait_result, _) = tokio::join!(wait_on_child, lines); - tracing::info!(%job_id, "child process '{child_name}' for {job_id} took {}ms, mem_peak: {:?}", start.elapsed().as_millis(), mem_peak); + tracing::info!(%job_id, "child process '{child_name}' for {worker_name}/{job_id} took {}ms, mem_peak: {:?}", start.elapsed().as_millis(), mem_peak); match wait_result { _ if *too_many_logs.borrow() => Err(Error::ExecutionErr(format!( "logs or result reached limit. (current max size: {MAX_RESULT_SIZE} characters)" diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index bedea13afb..96e21b46d7 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -16,7 +16,7 @@ use sqlx::{types::Json, Pool, Postgres}; use std::{ collections::HashMap, sync::{ - atomic::{AtomicUsize, Ordering}, + atomic::{AtomicBool, AtomicUsize, Ordering}, Arc, }, time::Duration, @@ -205,10 +205,10 @@ lazy_static::lazy_static! { .and_then(|x| x.parse::().ok()) .unwrap_or(true); - pub static ref KEEP_JOB_DIR: bool = std::env::var("KEEP_JOB_DIR") - .ok() - .and_then(|x| x.parse::().ok()) - .unwrap_or(false); + pub static ref KEEP_JOB_DIR: AtomicBool = AtomicBool::new(std::env::var("KEEP_JOB_DIR") + .ok() + .and_then(|x| x.parse::().ok()) + .unwrap_or(false)); pub static ref NO_PROXY: Option = std::env::var("no_proxy").ok().or(std::env::var("NO_PROXY").ok()); pub static ref HTTP_PROXY: Option = std::env::var("http_proxy").ok().or(std::env::var("HTTP_PROXY").ok()); @@ -1339,7 +1339,8 @@ pub async fn run_worker/', + storage: 'setting' + }, + { + label: 'Expose Debug Metrics', + key: 'expose_debug_metrics', + fieldType: 'boolean', + tooltip: 'Expose additional metrics (require metrics to be enabled)', + storage: 'setting' + } + ], Telemetry: [ { label: 'Disable telemetry',