diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index ab3f4a7ca0..f193b19089 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -19,21 +19,31 @@ use windmill_api::{ DEFAULT_BODY_LIMIT, IS_SECURE, OAUTH_CLIENTS, REQUEST_SIZE_LIMIT, SAML_METADATA, SCIM_TOKEN, }; use windmill_common::{ - error, flow_status::FlowStatusModule, global_settings::{ + error, + flow_status::FlowStatusModule, + global_settings::{ BASE_URL_SETTING, BUNFIG_INSTALL_SCOPES_SETTING, DEFAULT_TAGS_PER_WORKSPACE_SETTING, EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, PIP_INDEX_URL_SETTING, REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING, RETENTION_PERIOD_SECS_SETTING, SAML_METADATA_SETTING, SCIM_TOKEN_SETTING, - }, jobs::QueuedJob, oauth2::REQUIRE_PREEXISTING_USER_FOR_OAUTH, server::load_server_config, users::truncate_token, worker::{ + }, + jobs::QueuedJob, + oauth2::REQUIRE_PREEXISTING_USER_FOR_OAUTH, + server::load_server_config, + users::truncate_token, + worker::{ load_worker_config, reload_custom_tags_setting, DEFAULT_TAGS_PER_WORKSPACE, SERVER_CONFIG, WORKER_CONFIG, - }, BASE_URL, DB, METRICS_DEBUG_ENABLED, METRICS_ENABLED + }, + BASE_URL, DB, METRICS_DEBUG_ENABLED, METRICS_ENABLED, }; use windmill_queue::cancel_job; use windmill_worker::{ - create_token_for_owner, handle_job_error, AuthedClient, SameWorkerPayload, SendResult, BUNFIG_INSTALL_SCOPES, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, NPM_CONFIG_REGISTRY, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, SCRIPT_TOKEN_EXPIRY + create_token_for_owner, handle_job_error, AuthedClient, SameWorkerPayload, SendResult, + BUNFIG_INSTALL_SCOPES, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, NPM_CONFIG_REGISTRY, + PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, SCRIPT_TOKEN_EXPIRY, }; #[cfg(feature = "enterprise")] @@ -108,10 +118,8 @@ pub async fn initial_load( reload_worker_config(&db, tx, false).await; } - if server_mode { - if let Err(e) = reload_custom_tags_setting(db).await { - tracing::error!("Error reloading custom tags: {:?}", e) - } + if let Err(e) = reload_custom_tags_setting(db).await { + tracing::error!("Error reloading custom tags: {:?}", e) } if let Err(e) = reload_base_url_setting(db).await { @@ -603,7 +611,10 @@ pub async fn monitor_db( }; let expose_queue_metrics_f = async { - if !initial_load && METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) && server_mode { + if !initial_load + && METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed) + && server_mode + { expose_queue_metrics(&db).await; } }; @@ -793,8 +804,10 @@ async fn handle_zombie_jobs } } - let mut timeout_query = "SELECT * FROM queue WHERE last_ping < now() - ($1 || ' seconds')::interval - AND running = true AND job_kind NOT IN ('flow', 'flowpreview', 'singlescriptflow')".to_string(); + let mut timeout_query = + "SELECT * FROM queue WHERE last_ping < now() - ($1 || ' seconds')::interval + AND running = true AND job_kind NOT IN ('flow', 'flowpreview', 'singlescriptflow')" + .to_string(); if *RESTART_ZOMBIE_JOBS { timeout_query.push_str(" AND same_worker = true"); }; @@ -813,7 +826,8 @@ async fn handle_zombie_jobs tracing::info!("timedout zombie job {} {}", job.id, job.workspace_id,); // since the job is unrecoverable, the same worker queue should never be sent anything - let (same_worker_tx_never_used, _same_worker_rx_never_used) = mpsc::channel::(1); + let (same_worker_tx_never_used, _same_worker_rx_never_used) = + mpsc::channel::(1); let (send_result_never_used, _send_result_rx_never_used) = mpsc::channel::(1); let token = create_token_for_owner( @@ -860,12 +874,10 @@ async fn handle_zombie_jobs } } - async fn handle_zombie_flows( db: &DB, rsmq: Option, ) -> error::Result<()> { - let flows = sqlx::query_as::<_, QueuedJob>( r#" SELECT * @@ -877,12 +889,13 @@ async fn handle_zombie_flows( .fetch_all(db) .await?; - - for flow in flows { let status = flow.parse_flow_status(); - if status.is_some_and(|s| s.modules.get(0).is_some_and(|x| matches!(x, FlowStatusModule::WaitingForPriorSteps { .. }))) - { + if status.is_some_and(|s| { + s.modules + .get(0) + .is_some_and(|x| matches!(x, FlowStatusModule::WaitingForPriorSteps { .. })) + }) { tracing::error!( "Zombie flow detected: {} in workspace {}. It hasn't started yet, restarting it.", flow.id, @@ -897,7 +910,13 @@ async fn handle_zombie_flows( .await?; } else { let id = flow.id.clone(); - cancel_zombie_flow_job(db, flow, &rsmq, format!("Flow {} cancelled as it was hanging in between 2 steps", id)).await?; + cancel_zombie_flow_job( + db, + flow, + &rsmq, + format!("Flow {} cancelled as it was hanging in between 2 steps", id), + ) + .await?; } } @@ -907,15 +926,18 @@ async fn handle_zombie_flows( FROM parallel_monitor_lock WHERE last_ping IS NOT NULL AND last_ping < NOW() - ($1 || ' seconds')::interval RETURNING parent_flow_id, job_id, last_ping - ", FLOW_ZOMBIE_TRANSITION_TIMEOUT.as_str() + ", + FLOW_ZOMBIE_TRANSITION_TIMEOUT.as_str() ) .fetch_all(db) .await?; for flow in flows2 { - let in_queue = sqlx::query_as::<_, QueuedJob>( - "SELECT * FROM queue WHERE id = $1 AND running = true", - ).bind(flow.parent_flow_id).fetch_optional(db).await?; + let in_queue = + sqlx::query_as::<_, QueuedJob>("SELECT * FROM queue WHERE id = $1 AND running = true") + .bind(flow.parent_flow_id) + .fetch_optional(db) + .await?; if let Some(job) = in_queue { tracing::error!( "parallel Zombie flow detected: {} in workspace {}. Last ping was: {:?}.", @@ -931,9 +953,14 @@ async fn handle_zombie_flows( } } Ok(()) - } +} -async fn cancel_zombie_flow_job(db: &Pool, flow: QueuedJob, rsmq: &Option, message: String) -> Result<(), error::Error> { +async fn cancel_zombie_flow_job( + db: &Pool, + flow: QueuedJob, + rsmq: &Option, + message: String, +) -> Result<(), error::Error> { let tx = db.begin().await.unwrap(); tracing::error!( "zombie flow detected: {} in workspace {}. Cancelling it.", @@ -959,4 +986,4 @@ async fn cancel_zombie_flow_job(db: &Pool, flow: QueuedJob, rsmq: &Opt .await?; ntx.commit().await?; Ok(()) -} \ No newline at end of file +}