From 57dc62b298d73f658cd17a52e37ebebafce78c8b Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Wed, 25 Oct 2023 23:32:39 +0200 Subject: [PATCH] fix: prometheus metrics are an instance settings --- backend/src/main.rs | 33 +++++++++++++----- backend/src/monitor.rs | 34 ++++++++++++++++--- backend/windmill-api/src/webhook_util.rs | 4 +-- .../windmill-common/src/global_settings.rs | 1 + backend/windmill-common/src/lib.rs | 34 +++++++++++++++---- backend/windmill-queue/src/jobs.rs | 10 +++--- backend/windmill-worker/src/worker.rs | 14 ++++---- docker-compose.yml | 7 ---- .../lib/components/InstanceSettings.svelte | 8 +++++ 9 files changed, 103 insertions(+), 42 deletions(-) diff --git a/backend/src/main.rs b/backend/src/main.rs index 94b755adf1..6c851461a8 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -23,14 +23,14 @@ use tokio::{ use windmill_api::HTTP_CLIENT; use windmill_common::{ global_settings::{ - BASE_URL_SETTING, CUSTOM_TAGS_SETTING, DISABLE_STATS_SETTING, ENV_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, }, stats::schedule_stats, utils::rd_string, worker::{reload_custom_tags_setting, WORKER_GROUP}, - DB, METRICS_ADDR, + DB, METRICS_ADDR, METRICS_ENABLED, }; use windmill_worker::{ BUN_CACHE_DIR, BUN_TMP_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, @@ -65,6 +65,10 @@ pub enum Mode { async fn main() -> anyhow::Result<()> { dotenv::dotenv().ok(); + if std::env::var("RUST_LOG").is_err() { + std::env::set_var("RUST_LOG", "info") + } + #[cfg(not(feature = "flamegraph"))] windmill_common::tracing_init::initialize_tracing(); @@ -124,7 +128,6 @@ async fn main() -> anyhow::Result<()> { "We STRONGLY recommend using at most 1 worker per container, use at your own risks" ); } - let metrics_addr: Option = *METRICS_ADDR; let server_mode = !std::env::var("DISABLE_SERVER") .ok() @@ -338,15 +341,26 @@ Windmill Community Edition {GIT_VERSION} NPM_CONFIG_REGISTRY_SETTING => { reload_npm_config_registry_setting(&db).await }, - REQUEST_SIZE_LIMIT_SETTING => { - tracing::info!("Request limit size change detected, killing server expecting to be restarted"); - // we wait a bit randomly to avoid having all servers shutdown at same time + EXPOSE_METRICS => { + 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 => { + if server_mode { + tracing::info!("Request limit size change detected, killing server expecting to be restarted"); + // we wait a bit randomly to avoid having all servers 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"); + } + } + }, DISABLE_STATS_SETTING => {}, a @_ => { tracing::info!("Unrecognized Global Setting Change Payload: {:?}", a); @@ -381,12 +395,13 @@ Windmill Community Edition {GIT_VERSION} }; let metrics_f = async { - if let Some(_addr) = metrics_addr { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { #[cfg(not(feature = "enterprise"))] - panic!("Metrics are only available in the Enterprise Edition"); + tracing::error!("Metrics are only available in the EE, ignoring..."); #[cfg(feature = "enterprise")] - windmill_common::serve_metrics(_addr, rx.resubscribe(), num_workers > 0).await; + windmill_common::serve_metrics(*METRICS_ADDR, rx.resubscribe(), num_workers > 0) + .await; } Ok(()) as anyhow::Result<()> }; diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index ff95e9b1b3..72a623f928 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1,4 +1,11 @@ -use std::{collections::HashMap, fmt::Display, ops::Mul, str::FromStr, sync::Arc, time::Duration}; +use std::{ + collections::HashMap, + fmt::Display, + ops::Mul, + str::FromStr, + sync::{atomic::Ordering, Arc}, + time::Duration, +}; use serde::de::DeserializeOwned; use sqlx::{Pool, Postgres}; @@ -14,7 +21,7 @@ use windmill_api::{ use windmill_common::{ error, global_settings::{ - BASE_URL_SETTING, EXTRA_PIP_INDEX_URL_SETTING, LICENSE_KEY_SETTING, + BASE_URL_SETTING, EXPOSE_METRICS, EXTRA_PIP_INDEX_URL_SETTING, LICENSE_KEY_SETTING, NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, REQUEST_SIZE_LIMIT_SETTING, RETENTION_PERIOD_SECS_SETTING, }, @@ -76,6 +83,9 @@ pub async fn initial_load( worker_mode: bool, server_mode: bool, ) { + if let Err(e) = load_metrics_enabled(db).await { + tracing::error!("Error reloading loading metrics: {e}"); + } let reload_worker_config_f = async { if worker_mode { reload_worker_config(&db, tx, false).await; @@ -144,6 +154,20 @@ pub async fn initial_load( ); } +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 + ) + .fetch_optional(db) + .await; + match metrics_enabled { + Ok(Some(serde_json::Value::Bool(t))) => METRICS_ENABLED.store(t, Ordering::Relaxed), + _ => (), + }; + Ok(()) +} + pub async fn delete_expired_items(db: &DB) -> () { let tokens_deleted_r: std::result::Result, _> = sqlx::query_scalar( "DELETE FROM token WHERE expiration <= now() @@ -413,7 +437,7 @@ pub async fn monitor_db .ok() .unwrap_or_else(|| vec![]); - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_ZOMBIE_RESTART_COUNT.inc_by(restarted.len() as _); } for r in restarted { @@ -613,7 +637,7 @@ async fn handle_zombie_jobs .ok() .unwrap_or_else(|| vec![]); - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_ZOMBIE_DELETE_COUNT.inc_by(timeouts.len() as _); } diff --git a/backend/windmill-api/src/webhook_util.rs b/backend/windmill-api/src/webhook_util.rs index 1be01a1222..d3b6e412af 100644 --- a/backend/windmill-api/src/webhook_util.rs +++ b/backend/windmill-api/src/webhook_util.rs @@ -109,13 +109,13 @@ impl WebhookShared { } }; if let Some(url) = webhook_opt { - let timer = if *METRICS_ENABLED { Some(WEBHOOK_REQUEST_COUNT.start_timer()) } else { None }; + let timer = if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { Some(WEBHOOK_REQUEST_COUNT.start_timer()) } else { None }; let _ = client.post(url).json(&message).send().await; timer.map(|x| x.stop_and_record()); } }, Some(WebhookPayload::InstanceEvent(event)) => { - if *METRICS_ENABLED { Some(WEBHOOK_REQUEST_COUNT.start_timer()) } else { None }; + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { Some(WEBHOOK_REQUEST_COUNT.start_timer()) } else { None }; let r = client.post(INSTANCE_EVENTS_WEBHOOK.as_ref().unwrap()).json(&event).send().await; if let Err(e) = r { tracing::error!("Error sending instance event: {}", e); diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 239eb6936f..3ed0b0ba07 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -9,6 +9,7 @@ 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 ENV_SETTINGS: [&str; 54] = [ "DISABLE_NSJAIL", diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index b3bcb6c484..45d5e97fc7 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -6,7 +6,10 @@ * LICENSE-AGPL for a copy of the license. */ -use std::{net::SocketAddr, sync::Arc}; +use std::{ + net::SocketAddr, + sync::{atomic::AtomicBool, Arc}, +}; use error::Error; use scripts::ScriptLang; @@ -38,17 +41,24 @@ pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50; pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 5; lazy_static::lazy_static! { - pub static ref METRICS_ADDR: Option = std::env::var("METRICS_ADDR") + pub static ref METRICS_PORT: u16 = std::env::var("METRICS_PORT") + .ok() + .and_then(|s| s.parse::().ok()) + .unwrap_or(8001); + + pub static ref METRICS_ADDR: SocketAddr = std::env::var("METRICS_ADDR") .ok() .map(|s| { s.parse::() - .map(|b| b.then(|| SocketAddr::from(([0, 0, 0, 0], 8001)))) + .map(|b| b.then(|| SocketAddr::from(([0, 0, 0, 0], *METRICS_PORT)))) .or_else(|_| s.parse::().map(Some)) }) .transpose().ok() .flatten() - .flatten(); - pub static ref METRICS_ENABLED: bool = METRICS_ADDR.is_some(); + .flatten() + .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 BASE_URL: Arc> = Arc::new(RwLock::new("".to_string())); pub static ref IS_READY: std::sync::atomic::AtomicBool = std::sync::atomic::AtomicBool::new(false); } @@ -92,9 +102,14 @@ pub async fn serve_metrics( ) -> JoinHandle<()> { use std::sync::atomic::Ordering; - use axum::{routing::get, Router}; + use axum::{ + routing::{get, post}, + Router, + }; use hyper::StatusCode; - let router = Router::new().route("/metrics", get(metrics)); + let router = Router::new() + .route("/metrics", get(metrics)) + .route("/reset", post(reset)); let router = if ready_worker_endpoint { router.route( @@ -112,6 +127,7 @@ pub async fn serve_metrics( }; tokio::spawn(async move { + tracing::info!("Serving metrics at: {addr}"); if let Err(e) = axum::Server::bind(&addr) .serve(router.into_make_service()) .with_graceful_shutdown(async { @@ -132,6 +148,10 @@ async fn metrics() -> Result { .map_err(anyhow::Error::from)?) } +async fn reset() -> () { + todo!() +} + #[cfg(feature = "sqlx")] pub async fn connect_db(server_mode: bool) -> anyhow::Result> { use anyhow::Context; diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 7af25185ac..1cae723c27 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -195,7 +195,7 @@ pub async fn add_completed_job_error, rsmq: Option, ) -> Result { - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { metrics.map(|m| m.worker_execution_failed.inc()); } let result = WrappedError { error: e }; @@ -1076,7 +1076,7 @@ pub async fn pull( || pulled_job.concurrent_limit.is_none() || pulled_job.canceled { - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_PULL_COUNT.inc(); } return Ok(Option::Some(pulled_job)); @@ -1162,7 +1162,7 @@ pub async fn pull( concurrent_jobs_for_this_script ); if concurrent_jobs_for_this_script <= job_custom_concurrent_limit { - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_PULL_COUNT.inc(); } tx.commit().await?; @@ -1503,7 +1503,7 @@ pub async fn delete_job<'c, R: rsmq_async::RsmqConnection + Clone + Send>( w_id: &str, job_id: Uuid, ) -> windmill_common::error::Result> { - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_DELETE_COUNT.inc(); } let job_removed = sqlx::query_scalar!( @@ -2211,7 +2211,7 @@ pub async fn push<'c, T: Serialize + Send + Sync, R: rsmq_async::RsmqConnection .map_err(|e| Error::InternalErr(format!("Could not insert into queue {job_id}: {e}")))?; // TODO: technically the job isn't queued yet, as the transaction can be rolled back. Should be solved when moving these metrics to the queue abstraction. - if *METRICS_ENABLED { + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { QUEUE_PUSH_COUNT.inc(); } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 5efffc307d..bedea13afb 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -656,7 +656,7 @@ pub async fn run_worker 3 { last_checked_suspended = Instant::now(); true @@ -1188,7 +1188,7 @@ pub async fn run_worker, language: &Option, ) -> Option { - let metrics = if *METRICS_ENABLED { + let metrics = if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { Some(Metrics { worker_execution_failed: worker_execution_failed .get(language) diff --git a/docker-compose.yml b/docker-compose.yml index ac6fef381f..a89bc6645a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -30,10 +30,7 @@ services: - 8000 environment: - DATABASE_URL=${DATABASE_URL} - - RUST_LOG=info - MODE=server - ## You can set the number of workers to 1 and not need any separate worker service but not recommended - - METRICS_ADDR=false # (ee only, if set to true, metrics will be exposed on port 8001) depends_on: db: condition: service_healthy @@ -50,10 +47,8 @@ services: restart: unless-stopped environment: - DATABASE_URL=${DATABASE_URL} - - RUST_LOG=info - MODE=worker - KEEP_JOB_DIR=false - - METRICS_ADDR=false - WORKER_GROUP=default depends_on: db: @@ -78,10 +73,8 @@ services: restart: unless-stopped environment: - DATABASE_URL=${DATABASE_URL} - - RUST_LOG=info - MODE=worker - WORKER_GROUP=native - - METRICS_ADDR=false # (ee only, if set to true, metrics will be exposed on port 8001) depends_on: db: condition: service_healthy diff --git a/frontend/src/lib/components/InstanceSettings.svelte b/frontend/src/lib/components/InstanceSettings.svelte index 27ad8e16ef..1bbfddb178 100644 --- a/frontend/src/lib/components/InstanceSettings.svelte +++ b/frontend/src/lib/components/InstanceSettings.svelte @@ -84,6 +84,14 @@ storage: 'setting', ee_only: 'You can still set this setting by using NPM_CONFIG_REGISTRY as env variable to the worker containers' + }, + { + label: 'Expose metrics', + description: 'Expose prometheus metrics for workers and servers on port 8001 at /metrics', + key: 'expose_metrics', + fieldType: 'boolean', + storage: 'setting', + ee_only: 'No workaround around this' } ], SMTP: [