From 8f9bc6d3d3f4071bcecbe0891420fe2332c09a0a Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Fri, 27 Oct 2023 21:00:28 +0200 Subject: [PATCH] feat: refactor metrics and add performance debug metrics (#2520) * all * feat: add debug metrics --- backend/src/main.rs | 18 +- backend/src/monitor.rs | 32 ++ backend/windmill-worker/src/bun_executor.rs | 10 +- .../windmill-worker/src/dedicated_worker.rs | 8 +- .../windmill-worker/src/python_executor.rs | 8 +- backend/windmill-worker/src/worker.rs | 310 ++++++++++++++---- 6 files changed, 307 insertions(+), 79 deletions(-) diff --git a/backend/src/main.rs b/backend/src/main.rs index c1141246cd..74f9df37db 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -41,7 +41,7 @@ use windmill_worker::{ }; use crate::monitor::{ - initial_load, load_keep_job_dir, monitor_db, reload_base_url_setting, + initial_load, load_keep_job_dir, monitor_db, monitor_pool, 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, }; @@ -229,6 +229,8 @@ Windmill Community Edition {GIT_VERSION} monitor_db(&db, &base_internal_url, rsmq.clone(), server_mode).await; + monitor_pool(&db).await; + if std::env::var("BASE_INTERNAL_URL").is_ok() { tracing::warn!("BASE_INTERNAL_URL is now unecessary and ignored, you can remove it."); } @@ -346,12 +348,14 @@ Windmill Community Edition {GIT_VERSION} 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); - tokio::time::sleep(Duration::from_secs(rd_delay)).await; - if let Err(e) = tx.send(()) { - tracing::error!(error = %e, "Could not send killpill to server"); + if n.payload() != EXPOSE_DEBUG_METRICS_SETTING || 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); + tokio::time::sleep(Duration::from_secs(rd_delay)).await; + if let Err(e) = tx.send(()) { + tracing::error!(error = %e, "Could not send killpill to server"); + } } }, REQUEST_SIZE_LIMIT_SETTING => { diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index d49d48ff73..6fb1127ca5 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -414,6 +414,38 @@ pub async fn reload_setting( Ok(()) } +pub async fn monitor_pool(db: &DB) { + if METRICS_ENABLED.load(Ordering::Relaxed) { + let db = db.clone(); + tokio::spawn(async move { + let active_pool_connections: prometheus::IntGauge = prometheus::register_int_gauge!( + "pool_connections_active", + "Number of active postgresql connections in the pool" + ) + .unwrap(); + + let idle_pool_connections: prometheus::IntGauge = prometheus::register_int_gauge!( + "pool_connections_idle", + "Number of idle postgresql connections in the pool" + ) + .unwrap(); + + let max_pool_connections: prometheus::IntGauge = prometheus::register_int_gauge!( + "pool_connections_max", + "Number of max postgresql connections in the pool" + ) + .unwrap(); + + max_pool_connections.set(db.options().get_max_connections() as i64); + loop { + active_pool_connections.set(db.size() as i64); + idle_pool_connections.set(db.num_idle() as i64); + tokio::time::sleep(Duration::from_secs(30)).await; + } + }); + } +} + pub async fn monitor_db( db: &Pool, base_internal_url: &str, diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index a2feccec46..f9d6ffb12f 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -7,7 +7,7 @@ use serde_json::value::RawValue; use uuid::Uuid; #[cfg(feature = "enterprise")] -use crate::{common::build_envs_map, JobCompleted}; +use crate::common::build_envs_map; use crate::{ common::{ @@ -26,7 +26,7 @@ use tokio::{ use tokio::io::AsyncReadExt; #[cfg(feature = "enterprise")] -use tokio::sync::mpsc::{Receiver, Sender}; +use tokio::sync::mpsc::Receiver; #[cfg(feature = "enterprise")] use windmill_common::variables; @@ -526,6 +526,8 @@ pub async fn get_common_bun_proc_envs(base_internal_url: &str) -> HashMap, + job_completed_tx: JobCompletedSender, jobs_rx: Receiver>, killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> Result<()> { - use crate::dedicated_worker::handle_dedicated_process; - let mut logs = "".to_string(); let mut mem_peak: i32 = 0; let _ = write_file(job_dir, "main.ts", inner_content).await?; diff --git a/backend/windmill-worker/src/dedicated_worker.rs b/backend/windmill-worker/src/dedicated_worker.rs index 0571da6e73..e43248d279 100644 --- a/backend/windmill-worker/src/dedicated_worker.rs +++ b/backend/windmill-worker/src/dedicated_worker.rs @@ -14,12 +14,14 @@ use std::{collections::VecDeque, process::Stdio, sync::Arc}; use anyhow::Context; -use crate::{common::start_child_process, JobCompleted, MAX_BUFFERED_DEDICATED_JOBS}; +use crate::{ + common::start_child_process, JobCompleted, JobCompletedSender, MAX_BUFFERED_DEDICATED_JOBS, +}; use futures::{future, Future}; use std::{collections::HashMap, task::Poll}; -use tokio::sync::mpsc::{Receiver, Sender}; +use tokio::sync::mpsc::Receiver; fn conditional_polling( fut: impl Future, @@ -50,7 +52,7 @@ pub async fn handle_dedicated_process( common_bun_proc_envs: HashMap, args: Vec<&str>, mut killpill_rx: tokio::sync::broadcast::Receiver<()>, - job_completed_tx: Sender, + job_completed_tx: JobCompletedSender, token: &str, mut jobs_rx: Receiver>, ) -> std::result::Result<(), error::Error> { diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 2671227880..2f04b102e5 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -752,11 +752,11 @@ pub async fn handle_python_reqs( #[cfg(feature = "enterprise")] use std::sync::Arc; -#[cfg(feature = "enterprise")] -use tokio::sync::mpsc::Sender; #[cfg(feature = "enterprise")] -use crate::{common::build_envs_map, dedicated_worker::handle_dedicated_process, JobCompleted}; +use crate::JobCompletedSender; +#[cfg(feature = "enterprise")] +use crate::{common::build_envs_map, dedicated_worker::handle_dedicated_process}; #[cfg(feature = "enterprise")] use tokio::sync::mpsc::Receiver; #[cfg(feature = "enterprise")] @@ -774,7 +774,7 @@ pub async fn start_worker( w_id: &str, script_path: &str, token: &str, - job_completed_tx: Sender, + job_completed_tx: JobCompletedSender, jobs_rx: Receiver>, killpill_rx: tokio::sync::broadcast::Receiver<()>, ) -> error::Result<()> { diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index e79c0e61a9..787915e0cb 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -9,7 +9,10 @@ use anyhow::Result; use const_format::concatcp; use itertools::Itertools; -use prometheus::IntCounter; +use prometheus::{ + core::{AtomicI64, GenericGauge}, + IntCounter, +}; use reqwest::Response; use serde::{de::DeserializeOwned, Deserialize, Serialize}; use sqlx::{types::Json, Pool, Postgres}; @@ -33,7 +36,7 @@ use windmill_common::{ worker::{ to_raw_value, to_raw_value_owned, update_ping, CLOUD_HOSTED, WORKER_CONFIG, WORKER_GROUP, }, - DB, IS_READY, METRICS_ENABLED, + DB, IS_READY, METRICS_DEBUG_ENABLED, METRICS_ENABLED, }; use windmill_queue::{ canceled_job_to_result, empty_args, get_queued_job, pull, push, register_metric, PushArgs, @@ -238,11 +241,11 @@ lazy_static::lazy_static! { .and_then(|x| x.parse::().ok()) .unwrap_or(100); - static ref WORKER_STARTED: prometheus::IntGauge = prometheus::register_int_gauge!( + static ref WORKER_STARTED: Option = if METRICS_ENABLED.load(Ordering::Relaxed) { Some(prometheus::register_int_gauge!( "worker_started", "Total number of workers started." ) - .unwrap(); + .unwrap()) } else { None }; static ref WORKER_UPTIME_OPTS: prometheus::Opts = prometheus::opts!( "worker_uptime", @@ -455,6 +458,8 @@ async fn handle_receive_completed_job< same_worker_tx: Sender, rsmq: Option, worker_name: &str, + worker_save_completed_job_duration: Option>, + worker_flow_transition_duration: Option>, ) { let token = jc.token.clone(); let workspace = jc.job.workspace_id.clone(); @@ -474,6 +479,8 @@ async fn handle_receive_completed_job< same_worker_tx.clone(), rsmq.clone(), worker_name, + worker_save_completed_job_duration, + worker_flow_transition_duration, ) .await { @@ -493,6 +500,28 @@ async fn handle_receive_completed_job< } } +#[derive(Clone)] +pub struct JobCompletedSender( + Sender, + Option>>, + Option>, +); + +impl JobCompletedSender { + pub async fn send( + &self, + jc: JobCompleted, + ) -> Result<(), tokio::sync::mpsc::error::SendError> { + if let Some(wj) = self.1.as_ref() { + wj.inc() + } + let timer = self.2.as_ref().map(|x| x.start_timer()); + let r = self.0.send(jc).await; + timer.map(|x| x.stop_and_record()); + r + } +} + pub async fn run_worker( db: &Pool, worker_instance: &str, @@ -542,59 +571,165 @@ pub async fn run_worker = all_langs - // .clone() - // .into_iter() - // .map(|x| { - // ( - // x.clone(), - // prometheus::register_counter!(prometheus::Opts::new( - // "worker_execution_duration_counter", - // "Total number of seconds spent executing jobs" - // ) - // .const_label("name", &worker_name) - // .const_label("language", x.map(|x| x.as_str()).unwrap_or("none"))) - // .expect("register prometheus metric"), - // ) - // }) - // .collect(); + let worker_sleep_duration_counter = if METRICS_ENABLED.load(Ordering::Relaxed) { + Some( + prometheus::register_counter!(prometheus::opts!( + "worker_sleep_duration_counter", + "Total number of seconds spent sleeping between pulling jobs from the queue" + ) + .const_label("name", &worker_name)) + .expect("register prometheus metric"), + ) + } else { + None + }; - let worker_sleep_duration_counter = prometheus::register_counter!(prometheus::opts!( - "worker_sleep_duration_counter", - "Total number of seconds spent sleeping between pulling jobs from the queue" - ) - .const_label("name", &worker_name)) - .expect("register prometheus metric"); + let worker_pull_duration = if METRICS_ENABLED.load(Ordering::Relaxed) { + Some( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_pull_duration", + "Duration pulling next job", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + ) + } else { + None + }; - let worker_pull_duration = prometheus::register_histogram!(prometheus::HistogramOpts::new( - "worker_pull_duration", - "Duration pulling next job", - ) - .const_label("name", &worker_name),) - .expect("register prometheus metric"); + let worker_job_completed_channel_queue = if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) + && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_int_gauge!(prometheus::opts!( + "worker_job_completed_channel_queue_length", + "Queue length of the job completed channel queue", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; - let worker_pull_duration_counter = prometheus::register_counter!(prometheus::opts!( + let worker_completed_channel_queue_send_duration = + if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_completed_channel_queue_duration", + "Duration sending job to completed job channel", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; + + let worker_save_completed_job_duration = if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) + && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_save__duration", + "Duration sending job to completed job channel", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; + + let worker_code_execution_duration = if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) + && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_code_execution_duration", + "Duration of executing the job itself without the saving or flow transition", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; + + let worker_flow_initial_transition_duration = if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) + && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_flow_initial_transition_duration", + "Duration sending job to completed job channel", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; + + let worker_flow_transition_duration = if METRICS_DEBUG_ENABLED.load(Ordering::Relaxed) + && METRICS_ENABLED.load(Ordering::Relaxed) + { + Some(Arc::new( + prometheus::register_histogram!(prometheus::HistogramOpts::new( + "worker_flow_transition_duration", + "Duration of doing a flow transition after the job is completed", + ) + .const_label("name", &worker_name),) + .expect("register prometheus metric"), + )) + } else { + None + }; + + let worker_pull_duration_counter = if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) + { + Some( + prometheus::register_counter!(prometheus::opts!( "worker_pull_duration_counter", "Total number of seconds spent pulling jobs (if growing large the db is undersized)" ) - .const_label("name", &worker_name)) - .expect("register prometheus metric"); + .const_label("name", &worker_name)) + .expect("register prometheus metric"), + ) + } else { + None + }; - let worker_busy: prometheus::IntGauge = prometheus::register_int_gauge!(prometheus::Opts::new( - "worker_busy", - "Is the worker busy executing a job?", - ) - .const_label("name", &worker_name)) - .unwrap(); + let worker_busy: Option = + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { + Some( + prometheus::register_int_gauge!(prometheus::Opts::new( + "worker_busy", + "Is the worker busy executing a job?", + ) + .const_label("name", &worker_name)) + .unwrap(), + ) + } else { + None + }; let mut jobs_executed = 0; - if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { - WORKER_STARTED.inc(); + if let Some(ws) = WORKER_STARTED.as_ref() { + ws.inc(); } let (_copy_to_bucket_tx, mut copy_to_bucket_rx) = mpsc::channel::<()>(2); @@ -631,6 +766,12 @@ pub async fn run_worker(3); + let job_completed_tx = JobCompletedSender( + job_completed_tx, + worker_job_completed_channel_queue.clone(), + worker_completed_channel_queue_send_duration, + ); + let db2 = db.clone(); let base_internal_url2 = base_internal_url.to_string(); let same_worker_tx2 = same_worker_tx.clone(); @@ -705,9 +846,16 @@ pub async fn run_worker 0 { - tokio::time::sleep(Duration::from_millis(100)).await; + tokio::time::sleep(Duration::from_millis(50)).await; } tracing::info!("finished processing all completed jobs"); @@ -972,10 +1135,12 @@ pub async fn run_worker 3 { last_checked_suspended = Instant::now(); true @@ -1117,7 +1282,9 @@ pub async fn run_worker { @@ -1454,12 +1625,18 @@ pub async fn process_completed_job, rsmq: Option, worker_name: &str, + worker_save_completed_job_duration: Option>, + worker_flow_transition_duration: Option>, ) -> windmill_common::error::Result<()> { if success { // println!("bef completed job{:?}", SystemTime::now()); if let Some(cached_path) = cached_res_path { save_in_cache(db, &job, cached_path.to_string(), &result).await; } + + let timer = worker_save_completed_job_duration + .as_ref() + .map(|x| x.start_timer()); add_completed_job( db, &job, @@ -1471,8 +1648,13 @@ pub async fn process_completed_job( same_worker_tx: Sender, base_internal_url: &str, rsmq: Option, - job_completed_tx: Sender, + job_completed_tx: JobCompletedSender, + worker_flow_initial_transition_duration: Option>, + worker_code_execution_duration: Option>, ) -> windmill_common::error::Result<()> { if job.canceled { return Err(Error::JsonErr(canceled_job_to_result(&job))); @@ -1815,6 +2000,7 @@ async fn handle_queued_job( match job.job_kind { JobKind::FlowPreview | JobKind::Flow => { let args = job.get_args(); + let timer = worker_flow_initial_transition_duration.map(|x| x.start_timer()); handle_flow( &job, db, @@ -1826,6 +2012,7 @@ async fn handle_queued_job( worker_name, ) .await?; + timer.map(|x| x.stop_and_record()); } _ => { let mut logs = "".to_string(); @@ -1905,7 +2092,8 @@ async fn handle_queued_job( .map(|x| x.to_owned()) .unwrap_or_else(|| serde_json::from_str("{}").unwrap())), _ => { - handle_code_execution_job( + let timer = worker_code_execution_duration.map(|x| x.start_timer()); + let r = handle_code_execution_job( job.as_ref(), db, client, @@ -1916,7 +2104,9 @@ async fn handle_queued_job( base_internal_url, worker_name, ) - .await + .await; + timer.map(|x| x.stop_and_record()); + r } }; @@ -1944,7 +2134,7 @@ async fn process_result( job: Arc, result: error::Result>, job_dir: &str, - job_completed_tx: Sender, + job_completed_tx: JobCompletedSender, logs: String, mem_peak: i32, cached_res_path: Option,