From f451748faaf42bff66fc34c9bd43a30f40375960 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 20 May 2024 18:44:03 +0200 Subject: [PATCH] allow dynamic heap profiling using jemalloc live prof_active --- Dockerfile | 2 ++ backend/Cargo.lock | 13 ++++++++ backend/Cargo.toml | 6 ++-- backend/src/main.rs | 29 ++++++++++++------ backend/src/monitor.rs | 69 +++++++++++++++++++++++++++++++++++++++++- 5 files changed, 107 insertions(+), 12 deletions(-) diff --git a/Dockerfile b/Dockerfile index fc9a64048c..c7cd51db68 100644 --- a/Dockerfile +++ b/Dockerfile @@ -185,6 +185,8 @@ RUN ln -s ${APP}/windmill /usr/local/bin/windmill RUN windmill cache +ENV _RJEM_MALLOC_CONF=prof:true,prof_active:false,lg_prof_interval:30,lg_prof_sample:21,prof_prefix:/tmp/jeprof + EXPOSE 8000 CMD ["windmill"] diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 5bc51817db..da72a5d0ab 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -8467,6 +8467,17 @@ dependencies = [ "digest 0.9.0", ] +[[package]] +name = "tikv-jemalloc-ctl" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "619bfed27d807b54f7f776b9430d4f8060e66ee138a28632ca898584d462c31c" +dependencies = [ + "libc", + "paste", + "tikv-jemalloc-sys", +] + [[package]] name = "tikv-jemalloc-sys" version = "0.5.4+5.3.0-patched" @@ -9710,6 +9721,8 @@ dependencies = [ "serde_json", "sha2 0.10.8", "sqlx", + "tikv-jemalloc-ctl", + "tikv-jemalloc-sys", "tikv-jemallocator", "tokio", "tokio-metrics", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index c6b26472c8..60e712af10 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -51,7 +51,7 @@ prometheus = ["windmill-common/prometheus", "windmill-api/prometheus", "windmill flow_testing = ["windmill-worker/flow_testing"] openidconnect = ["windmill-api/openidconnect"] cloud = ["windmill-queue/cloud"] -jemalloc = ["dep:tikv-jemallocator"] +jemalloc = ["dep:tikv-jemallocator", "dep:tikv-jemalloc-sys", "dep:tikv-jemalloc-ctl"] [dependencies] anyhow.workspace = true @@ -84,7 +84,9 @@ deno_core.workspace = true pg-embed = {git = "https://github.com/faokunega/pg-embed", optional = true, default-features = false, features = ['rt_tokio']} [target.'cfg(not(target_env = "msvc"))'.dependencies] -tikv-jemallocator = { optional = true, version = "0.5" } +tikv-jemallocator = { optional = true, version = "0.5", features = ["profiling"] } +tikv-jemalloc-sys = { optional = true, version = "^0.5", features = ["profiling"] } +tikv-jemalloc-ctl = { optional = true, version = "^0.5" } [dev-dependencies] serde_json.workspace = true diff --git a/backend/src/main.rs b/backend/src/main.rs index a5846e51b7..4e49143a6d 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -32,6 +32,9 @@ use windmill_common::{ DB, METRICS_ENABLED, }; +#[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] +use monitor::monitor_mem; + #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] use tikv_jemallocator::Jemalloc; @@ -52,13 +55,13 @@ use windmill_worker::{ }; use crate::monitor::{ - initial_load, load_keep_job_dir, load_require_preexisting_user, load_tag_per_workspace_enabled, - monitor_db, monitor_pool, reload_base_url_setting, reload_bunfig_install_scopes_setting, - reload_critical_error_channels_setting, reload_extra_pip_index_url_setting, - reload_hub_base_url_setting, reload_job_default_timeout_setting, reload_license_key, - reload_npm_config_registry_setting, reload_pip_index_url_setting, - reload_retention_period_setting, reload_scim_token_setting, reload_server_config, - reload_worker_config, + initial_load, load_keep_job_dir, load_metrics_debug_enabled, load_require_preexisting_user, + load_tag_per_workspace_enabled, monitor_db, monitor_pool, reload_base_url_setting, + reload_bunfig_install_scopes_setting, reload_critical_error_channels_setting, + reload_extra_pip_index_url_setting, reload_hub_base_url_setting, + reload_job_default_timeout_setting, reload_license_key, reload_npm_config_registry_setting, + reload_pip_index_url_setting, reload_retention_period_setting, reload_scim_token_setting, + reload_server_config, reload_worker_config, }; #[cfg(feature = "parquet")] @@ -321,6 +324,9 @@ Windmill Community Edition {GIT_VERSION} monitor_pool(&db).await; + #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] + monitor_mem().await; + let addr = SocketAddr::from((server_bind_address, port)); let rsmq2 = rsmq.clone(); @@ -481,8 +487,8 @@ Windmill Community Edition {GIT_VERSION} REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING => { load_require_preexisting_user(&db).await; }, - EXPOSE_METRICS_SETTING | EXPOSE_DEBUG_METRICS_SETTING => { - if n.payload() != EXPOSE_DEBUG_METRICS_SETTING || worker_mode { + EXPOSE_METRICS_SETTING => { + if worker_mode { 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); @@ -492,6 +498,11 @@ Windmill Community Edition {GIT_VERSION} } } }, + EXPOSE_DEBUG_METRICS_SETTING => { + if let Err(e) = load_metrics_debug_enabled(&db).await { + tracing::error!(error = %e, "Could not reload debug metrics setting"); + } + }, REQUEST_SIZE_LIMIT_SETTING => { if server_mode { tracing::info!("Request limit size change detected, killing server expecting to be restarted"); diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index b907b49e6a..60be1e837d 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -201,12 +201,79 @@ pub async fn load_tag_per_workspace_enabled(db: &DB) -> error::Result<()> { pub async fn load_metrics_debug_enabled(db: &DB) -> error::Result<()> { let metrics_enabled = load_value_from_global_settings(db, EXPOSE_DEBUG_METRICS_SETTING).await; match metrics_enabled { - Ok(Some(serde_json::Value::Bool(t))) => METRICS_DEBUG_ENABLED.store(t, Ordering::Relaxed), + Ok(Some(serde_json::Value::Bool(t))) => { + METRICS_DEBUG_ENABLED.store(t, Ordering::Relaxed); + #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] + if let Err(e) = set_prof_active(t) { + tracing::error!("Error setting jemalloc prof_active: {e:?}"); + } + }, _ => (), }; Ok(()) } +#[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] +#[derive(Debug, Clone)] +pub struct MallctlError { pub code: i32 } + +#[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] +fn set_prof_active(new_value: bool) -> Result<(), MallctlError> { + let option_name = std::ffi::CString::new("prof.active").unwrap(); + + tracing::info!("Setting jemalloc prof_active to {}", new_value); + let result = unsafe { + + tikv_jemalloc_sys::mallctl( + option_name.as_ptr(), // const char *name + std::ptr::null_mut(), // void *oldp + std::ptr::null_mut(), // size_t *oldlenp + &new_value as *const _ as *mut _, // void *newp + std::mem::size_of_val(&new_value) // size_t newlen + ) + }; + + if result != 0 { + return Err(MallctlError { code: result }); + } + + Ok(()) +} + +fn bytes_to_mb(bytes: u64) -> f64 { + const BYTES_PER_MB: f64 = 1_048_576.0; + bytes as f64 / BYTES_PER_MB +} + + +#[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] +pub async fn monitor_mem() { + + use std::time::Duration; + use tikv_jemalloc_ctl::{stats, epoch}; + + tokio::spawn(async move { + + // Obtain a MIB for the `epoch`, `stats.allocated`, and + // `atats.resident` keys: + let e = epoch::mib().unwrap(); + let allocated = stats::allocated::mib().unwrap(); + let resident = stats::resident::mib().unwrap(); + + loop { + // Many statistics are cached and only updated + // when the epoch is advanced: + e.advance().unwrap(); + + // Read statistics using MIB key: + let allocated = allocated.read().unwrap(); + let resident = resident.read().unwrap(); + tracing::info!("{} mb allocated/{} mb resident", bytes_to_mb(allocated as u64), bytes_to_mb(resident as u64)); + tokio::time::sleep(Duration::from_secs(10)).await; + } + }); +} + pub async fn load_keep_job_dir(db: &DB) { let value = load_value_from_global_settings(db, KEEP_JOB_DIR_SETTING).await; match value {