fix: move keep job directories and expose debug metrics to instance settings UI

This commit is contained in:
Ruben Fiszel
2023-10-26 13:08:39 +02:00
parent ea28163865
commit 55ceca1913
10 changed files with 117 additions and 82 deletions
+11 -7
View File
@@ -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);
+73 -63
View File
@@ -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<Vec<String>, _> = sqlx::query_scalar(
"DELETE FROM token WHERE expiration <= now()
@@ -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",
+2
View File
@@ -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<RwLock<String>> = Arc::new(RwLock::new("".to_string()));
pub static ref IS_READY: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false);
}
+1 -1
View File
@@ -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,
+2 -2
View File
@@ -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)"
+7 -6
View File
@@ -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::<bool>().ok())
.unwrap_or(true);
pub static ref KEEP_JOB_DIR: bool = std::env::var("KEEP_JOB_DIR")
.ok()
.and_then(|x| x.parse::<bool>().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::<bool>().ok())
.unwrap_or(false));
pub static ref NO_PROXY: Option<String> = std::env::var("no_proxy").ok().or(std::env::var("NO_PROXY").ok());
pub static ref HTTP_PROXY: Option<String> = std::env::var("http_proxy").ok().or(std::env::var("HTTP_PROXY").ok());
@@ -1339,7 +1339,8 @@ pub async fn run_worker<R: rsmq_async::RsmqConnection + Send + Sync + Clone + 's
.expect("no timer found")
.inc_by(duration);
if !*KEEP_JOB_DIR && !(arc_job.is_flow() && same_worker) {
if !KEEP_JOB_DIR.load(Ordering::Relaxed) && !(arc_job.is_flow() && same_worker)
{
let _ = tokio::fs::remove_dir_all(job_dir).await;
}
}
+2 -1
View File
@@ -8,6 +8,7 @@
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::time::Duration;
use crate::common::{hash_args, save_in_cache};
@@ -715,7 +716,7 @@ pub async fn update_flow_status_after_job_completion_internal<
};
if done {
if flow_job.same_worker && !*KEEP_JOB_DIR {
if flow_job.same_worker && !KEEP_JOB_DIR.load(Ordering::Relaxed) {
let _ = tokio::fs::remove_dir_all(format!("{worker_dir}/{}", flow_job.id)).await;
}
-1
View File
@@ -48,7 +48,6 @@ services:
environment:
- DATABASE_URL=${DATABASE_URL}
- MODE=worker
- KEEP_JOB_DIR=false
- WORKER_GROUP=default
depends_on:
db:
@@ -137,6 +137,22 @@
}
],
'SSO/OAuth': [],
Debug: [
{
label: 'Keep Job Directories',
key: 'keep_job_dir',
fieldType: 'boolean',
tooltip: 'Keep Job directories after execution at /tmp/windmill/<worker>/<job_id>',
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',