diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 08aa850e6d..d49d48ff73 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -422,7 +422,7 @@ pub async fn monitor_db db: &Pool, base_internal_url: &str, rsmq: Option, + worker_name: &str, ) { if *RESTART_ZOMBIE_JOBS { let restarted = sqlx::query!( @@ -688,11 +689,11 @@ async fn handle_zombie_jobs .unwrap_or_else(|| "no ping".to_string()), *ZOMBIE_JOB_TIMEOUT )), - None, true, same_worker_tx_never_used, "", rsmq.clone(), + worker_name, ) .await; } diff --git a/backend/windmill-common/src/jobs.rs b/backend/windmill-common/src/jobs.rs index 4257111f4d..e53e00c52e 100644 --- a/backend/windmill-common/src/jobs.rs +++ b/backend/windmill-common/src/jobs.rs @@ -330,8 +330,3 @@ pub async fn get_payload_tag_from_prefixed_path( }; Ok((payload, tag)) } - -#[derive(Clone)] -pub struct Metrics { - pub worker_execution_failed: prometheus::IntCounter, -} diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 32193850c5..ff064a84b7 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -6,7 +6,7 @@ * LICENSE-AGPL for a copy of the license. */ -use std::{collections::HashMap, vec}; +use std::{collections::HashMap, sync::Arc, vec}; use anyhow::Context; use async_recursion::async_recursion; @@ -19,6 +19,7 @@ use axum::{ }; use bigdecimal::ToPrimitive; use chrono::{DateTime, Duration, Utc}; +use prometheus::IntCounter; use reqwest::{ header::{HeaderMap, CONTENT_TYPE}, Client, StatusCode, @@ -29,6 +30,7 @@ use serde_json::{json, value::RawValue}; use sqlx::{types::Json, FromRow, Pool, Postgres, Transaction}; #[cfg(feature = "benchmark")] use std::time::Instant; +use tokio::sync::RwLock; use tracing::{instrument, Instrument}; use ulid::Ulid; use uuid::Uuid; @@ -43,8 +45,8 @@ use windmill_common::{ }, flows::{FlowModule, FlowModuleValue, FlowValue}, jobs::{ - get_payload_tag_from_prefixed_path, script_path_to_payload, JobKind, JobPayload, Metrics, - QueuedJob, RawCode, + get_payload_tag_from_prefixed_path, script_path_to_payload, JobKind, JobPayload, QueuedJob, + RawCode, }, schedule::{schedule_to_user, Schedule}, scripts::{ScriptHash, ScriptLang}, @@ -72,17 +74,22 @@ lazy_static::lazy_static! { "Total number of jobs pushed to the queue." ) .unwrap(); + static ref QUEUE_DELETE_COUNT: prometheus::IntCounter = prometheus::register_int_counter!( "queue_delete_count", "Total number of jobs deleted from the queue." ) .unwrap(); + static ref QUEUE_PULL_COUNT: prometheus::IntCounter = prometheus::register_int_counter!( "queue_pull_count", "Total number of jobs pulled from the queue." ) .unwrap(); + pub static ref WORKER_EXECUTION_FAILED: Arc>> = Arc::new(RwLock::new(HashMap::new())); + + } #[cfg(feature = "enterprise")] @@ -135,8 +142,8 @@ pub async fn cancel_job<'c: 'async_recursion>( format!("canceled by {username}: (force cancel: {force_cancel})"), job_running.mem_peak.unwrap_or(0), e, - None, rsmq.clone(), + "server", ) .await; if let Err(e) = add_job { @@ -185,6 +192,36 @@ pub struct WrappedError { pub error: serde_json::Value, } +pub async fn register_metric( + l: &Arc>>, + s: &str, + fnew: F2, + f: F, +) -> Option +where + F: FnOnce(&T) -> R, + F2: FnOnce(&str) -> (T, R), +{ + if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { + let lock = l.read().await; + + let counter = lock.get(s); + if let Some(counter) = counter { + let r = f(counter); + drop(lock); + Some(r) + } else { + drop(lock); + let (metric, r) = fnew(s); + let mut m = l.write().await; + (*m).insert(s.to_string(), metric); + Some(r) + } + } else { + None + } +} + #[instrument(level = "trace", skip_all)] pub async fn add_completed_job_error( db: &Pool, @@ -192,12 +229,27 @@ pub async fn add_completed_job_error, rsmq: Option, + worker_name: &str, ) -> Result { - if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { - metrics.map(|m| m.worker_execution_failed.inc()); - } + register_metric( + &WORKER_EXECUTION_FAILED, + &queued_job.tag, + |s| { + let counter = prometheus::register_int_counter!(prometheus::Opts::new( + "worker_execution_count", + "Number of executed jobs" + ) + .const_label("name", worker_name) + .const_label("tag", s)) + .expect("register prometheus metric"); + counter.inc(); + (counter, ()) + }, + |c| c.inc(), + ) + .await; + let result = WrappedError { error: e }; tracing::error!( "job {} did not succeed: {}", diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 96e21b46d7..e79c0e61a9 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -9,7 +9,7 @@ use anyhow::Result; use const_format::concatcp; use itertools::Itertools; -use prometheus::core::{AtomicU64, GenericCounter}; +use prometheus::IntCounter; use reqwest::Response; use serde::{de::DeserializeOwned, Deserialize, Serialize}; use sqlx::{types::Json, Pool, Postgres}; @@ -26,7 +26,7 @@ use uuid::Uuid; use windmill_common::{ error::{self, to_anyhow, Error}, flows::{FlowModule, FlowModuleValue, FlowValue}, - jobs::{JobKind, Metrics, QueuedJob}, + jobs::{JobKind, QueuedJob}, scripts::{get_full_hub_script_by_path, ScriptHash, ScriptLang}, users::SUPERADMIN_SECRET_EMAIL, utils::{rd_string, StripPath}, @@ -36,8 +36,8 @@ use windmill_common::{ DB, IS_READY, METRICS_ENABLED, }; use windmill_queue::{ - canceled_job_to_result, empty_args, get_queued_job, pull, push, PushArgs, PushIsolationLevel, - WrappedError, HTTP_CLIENT, + canceled_job_to_result, empty_args, get_queued_job, pull, push, register_metric, PushArgs, + PushIsolationLevel, WrappedError, HTTP_CLIENT, }; use serde_json::{json, value::RawValue, Value}; @@ -283,6 +283,10 @@ lazy_static::lazy_static! { pub static ref CAN_PULL: Arc> = Arc::new(RwLock::new(())); + pub static ref WORKER_EXECUTION_COUNT: Arc>> = Arc::new(RwLock::new(HashMap::new())); + pub static ref WORKER_EXECUTION_DURATION_COUNTER: Arc>> = Arc::new(RwLock::new(HashMap::new())); + + pub static ref WORKER_EXECUTION_DURATION: Arc>> = Arc::new(RwLock::new(HashMap::new())); } //only matter if CLOUD_HOSTED @@ -445,14 +449,13 @@ async fn handle_receive_completed_job< R: rsmq_async::RsmqConnection + Send + Sync + Clone + 'static, >( jc: JobCompleted, - worker_execution_failed: HashMap, GenericCounter>, base_internal_url: String, db: Pool, worker_dir: String, same_worker_tx: Sender, rsmq: Option, + worker_name: &str, ) { - let metrics = build_language_metrics(&worker_execution_failed.clone(), &jc.job.language); let token = jc.token.clone(); let workspace = jc.job.workspace_id.clone(); let client = AuthedClient { @@ -468,9 +471,9 @@ async fn handle_receive_completed_job< &client, &db, &worker_dir, - metrics.clone(), same_worker_tx.clone(), rsmq.clone(), + worker_name, ) .await { @@ -480,11 +483,11 @@ async fn handle_receive_completed_job< job.as_ref(), mem_peak, err, - metrics, false, same_worker_tx.clone(), &worker_dir, rsmq.clone(), + worker_name, ) .await; } @@ -543,55 +546,22 @@ pub async fn run_worker = all_langs - .clone() - .into_iter() - .map(|x| { - ( - x.clone(), - prometheus::register_histogram!(prometheus::HistogramOpts::new( - "worker_execution_duration", - "Duration between receiving a job and completing it", - ) - .const_label("name", &worker_name) - .const_label("language", x.map(|x| x.as_str()).unwrap_or("none"))) - .expect("register prometheus metric"), - ) - }) - .collect(); - - let worker_execution_duration_counter: HashMap<_, _> = 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_execution_duration_counter: HashMap<_, _> = 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 = prometheus::register_counter!(prometheus::opts!( "worker_sleep_duration_counter", @@ -614,39 +584,6 @@ pub async fn run_worker = all_langs - .clone() - .into_iter() - .map(|x| { - ( - x.clone(), - prometheus::register_int_counter!(prometheus::Opts::new( - "worker_execution_failed", - "Number of failed 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_execution_count: HashMap<_, _> = all_langs - .into_iter() - .map(|x| { - ( - x.clone(), - prometheus::register_int_counter!(prometheus::Opts::new( - "worker_execution_count", - "Number of executed 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_busy: prometheus::IntGauge = prometheus::register_int_gauge!(prometheus::Opts::new( "worker_busy", "Is the worker busy executing a job?", @@ -699,7 +636,6 @@ pub async fn run_worker, same_worker_tx: Sender, rsmq: Option, + worker_name: &str, ) -> windmill_common::error::Result<()> { if success { // println!("bef completed job{:?}", SystemTime::now()); @@ -1507,12 +1481,12 @@ pub async fn process_completed_job, - prometheus::core::GenericCounter, - >, - language: &Option, -) -> Option { - let metrics = if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { - Some(Metrics { - worker_execution_failed: worker_execution_failed - .get(language) - .expect("no timer found") - .clone(), - }) - } else { - None - }; - metrics -} +// fn build_language_metrics( +// worker_execution_failed: &HashMap< +// Option, +// prometheus::core::GenericCounter, +// >, +// language: &Option, +// ) -> Option { +// let metrics = if METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) { +// Some(Metrics { +// worker_execution_failed: worker_execution_failed +// .get(language) +// .expect("no timer found") +// .clone(), +// }) +// } else { +// None +// }; +// metrics +// } // pub async fn create_barrier_for_all_workers(num_workers: u32, sync_barrier: Arc>>) { // tracing::debug!("acquiring write lock"); @@ -1595,11 +1569,11 @@ pub async fn handle_job_error, unrecoverable: bool, same_worker_tx: Sender, worker_dir: &str, rsmq: Option, + worker_name: &str, ) { let err = match err { Error::JsonErr(err) => err, @@ -1614,8 +1588,8 @@ pub async fn handle_job_error( same_worker_tx, worker_dir, rsmq, + worker_name, ) .await?; } diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index b050718efa..6e1e800123 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -7,8 +7,8 @@ */ use std::collections::HashMap; -use std::sync::Arc; use std::sync::atomic::Ordering; +use std::sync::Arc; use std::time::Duration; use crate::common::{hash_args, save_in_cache}; @@ -28,7 +28,7 @@ use windmill_common::flow_status::{ ApprovalConditions, FlowStatusModuleWParent, Iterator, JobResult, }; use windmill_common::jobs::{ - script_hash_to_tag_and_limits, script_path_to_payload, JobPayload, Metrics, QueuedJob, RawCode, + script_hash_to_tag_and_limits, script_path_to_payload, JobPayload, QueuedJob, RawCode, }; use windmill_common::worker::to_raw_value; use windmill_common::{ @@ -60,12 +60,12 @@ pub async fn update_flow_status_after_job_completion< w_id: &str, success: bool, result: &'a RawValue, - metrics: Option, unrecoverable: bool, same_worker_tx: Sender, worker_dir: &str, stop_early_override: Option, rsmq: Option, + worker_name: &str, ) -> error::Result<()> { // this is manual tailrecursion because async_recursion blows up the stack let mut rec = update_flow_status_after_job_completion_internal( @@ -76,13 +76,13 @@ pub async fn update_flow_status_after_job_completion< w_id, success, result, - metrics.clone(), unrecoverable, same_worker_tx.clone(), worker_dir, stop_early_override, false, rsmq.clone(), + worker_name, ) .await?; while let Some(nrec) = rec { @@ -94,13 +94,13 @@ pub async fn update_flow_status_after_job_completion< w_id, nrec.success, nrec.result.as_ref(), - metrics.clone(), false, same_worker_tx.clone(), worker_dir, nrec.stop_early_override, nrec.skip_error_handler, rsmq.clone(), + worker_name, ) .await { @@ -114,13 +114,13 @@ pub async fn update_flow_status_after_job_completion< w_id, false, &to_raw_value(&Json(&WrappedError { error: json!(e.to_string()) })), - metrics.clone(), true, same_worker_tx.clone(), worker_dir, nrec.stop_early_override, nrec.skip_error_handler, rsmq.clone(), + worker_name, ) .await? } @@ -162,13 +162,13 @@ pub async fn update_flow_status_after_job_completion_internal< w_id: &str, mut success: bool, result: &'a RawValue, - metrics: Option, unrecoverable: bool, same_worker_tx: Sender, worker_dir: &str, stop_early_override: Option, skip_error_handler: bool, rsmq: Option, + worker_name: &str, ) -> error::Result> { let (should_continue_flow, flow_job, stop_early, skip_if_stop_early, nresult, is_failure_step) = { // tracing::debug!("UPDATE FLOW STATUS: {flow:?} {success} {result:?} {w_id} {depth}"); @@ -639,8 +639,8 @@ pub async fn update_flow_status_after_job_completion_internal< logs, 0, canceled_job_to_result(&flow_job), - metrics.clone(), rsmq.clone(), + worker_name, ) .await?; } else { @@ -694,6 +694,7 @@ pub async fn update_flow_status_after_job_completion_internal< same_worker_tx.clone(), worker_dir, rsmq.clone(), + worker_name, ) .await { @@ -705,8 +706,8 @@ pub async fn update_flow_status_after_job_completion_internal< "Unexpected error during flow chaining:\n".to_string(), 0, e, - metrics.clone(), rsmq.clone(), + worker_name, ) .await; true @@ -987,6 +988,7 @@ pub async fn handle_flow( same_worker_tx: Sender, worker_dir: &str, rsmq: Option, + worker_name: &str, ) -> anyhow::Result<()> { let value = flow_job .raw_flow @@ -1011,6 +1013,7 @@ pub async fn handle_flow( same_worker_tx, worker_dir, rsmq, + worker_name, ) .await?; Ok(()) @@ -1055,6 +1058,7 @@ async fn push_next_flow_job same_worker_tx: Sender, worker_dir: &str, rsmq: Option, + worker_name: &str, ) -> error::Result<()> { let job_root = flow_job .root_job @@ -1090,12 +1094,12 @@ async fn push_next_flow_job // it has to be an empty for loop event serde_json::from_str("[]").unwrap() }, - None, true, same_worker_tx, worker_dir, None, rsmq, + worker_name, ) .await; } @@ -1123,12 +1127,12 @@ async fn push_next_flow_job flow_job.workspace_id.as_str(), true, serde_json::from_str("\"stopped early\"").unwrap(), - None, true, same_worker_tx, worker_dir, Some(true), rsmq, + worker_name, ) .await; }