/* * Author: Ruben Fiszel * Copyright: Windmill Labs, Inc 2022 * This file and its contents are licensed under the AGPLv3 License. * Please see the included NOTICE for copyright information and * LICENSE-AGPL for a copy of the license. */ use anyhow::Context; use monitor::{ load_base_url, load_otel, reload_critical_alerts_on_db_oversize, reload_delete_logs_periodically_setting, reload_indexer_config, reload_instance_python_version_setting, reload_maven_repos_setting, reload_maven_settings_xml_setting, reload_no_default_maven_setting, reload_nuget_config_setting, reload_powershell_repo_pat_setting, reload_powershell_repo_url_setting, reload_ruby_repos_setting, reload_timeout_wait_result_setting, reload_workspace_registries_setting, send_current_log_file_to_object_store, send_logs_to_object_store, WORKERS_NAMES, }; use rand::Rng; use sqlx::{Pool, Postgres}; use std::{ collections::HashMap, fs::{create_dir_all, DirBuilder}, net::{IpAddr, Ipv4Addr, SocketAddr}, time::{Duration, Instant}, }; use strum::IntoEnumIterator; use tokio::{fs::File, io::AsyncReadExt, task::JoinHandle}; use uuid::Uuid; use windmill_api::HTTP_CLIENT; #[cfg(feature = "enterprise")] use windmill_common::ee_oss::{ maybe_renew_license_key_on_start, LICENSE_KEY_ID, LICENSE_KEY_VALID, }; use windmill_ai::ai_cache::bump_instance_ai_config_revision; use windmill_common::{ agent_workers::AgentConfig, global_settings::{ AI_CONFIG_SETTING, APP_WORKSPACED_ROUTE_SETTING, AUDIT_LOG_RETENTION_DAYS_SETTING, BASE_URL_SETTING, BUNFIG_INSTALL_SCOPES_SETTING, BUN_INSTALL_MIN_RELEASE_AGE_SETTING, CONCURRENCY_KEY_MAX_QUEUED_SETTING, CRITICAL_ALERTS_ON_DB_OVERSIZE_SETTING, CRITICAL_ALERTS_ON_TOKEN_EXPIRY_SETTING, CRITICAL_ALERT_MUTE_UI_SETTING, CRITICAL_ERROR_CHANNELS_SETTING, CUSTOM_TAGS_SETTING, DEFAULT_TAGS_PER_WORKSPACE_SETTING, DEFAULT_TAGS_WORKSPACES_SETTING, DISABLE_PASSWORD_LOGIN_SETTING, EMAIL_DOMAIN_SETTING, ENV_SETTINGS, EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING, FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING, HTTP_ROUTE_WORKSPACED_ROUTE_SETTING, HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING, INDEXER_SETTING, INSTANCE_EVENTS_WEBHOOK_SETTING, INSTANCE_PYTHON_VERSION_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING, JOB_ISOLATION_SETTING, JWT_SECRET_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MAVEN_REPOS_SETTING, MAVEN_SETTINGS_XML_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NO_DEFAULT_MAVEN_SETTING, NPM_CONFIG_REGISTRY_SETTING, NSJAIL_TMPFS_SIZE_MB_SETTING, NSJAIL_TMP_BACKING_SETTING, NUGET_CONFIG_SETTING, OAUTH_SETTING, OTEL_SETTING, OTEL_TRACING_PROXY_SETTING, PIP_INDEX_URL_SETTING, POWERSHELL_REPO_PAT_SETTING, POWERSHELL_REPO_URL_SETTING, PREVIEW_TAGS_OVERRIDE_SETTING, REQUEST_SIZE_LIMIT_SETTING, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING, RESTART_COORDINATION_SETTING, RETENTION_PERIOD_SECS_OVERRIDES_SETTING, RETENTION_PERIOD_SECS_SETTING, RUBY_REPOS_SETTING, SAML_METADATA_SETTING, SANDBOX_IMAGE_CACHE_MAX_MB_SETTING, SANDBOX_IMAGE_DEFAULT_REGISTRY_SETTING, SANDBOX_IMAGE_MAX_SIZE_MB_SETTING, SANDBOX_IMAGE_PULL_POLICY_SETTING, SANDBOX_REGISTRY_AUTH_SETTING, SCIM_TOKEN_SETTING, SMTP_SETTING, STORE_AUDIT_LOGS_S3_SETTING, TEAMS_SETTING, TIMEOUT_WAIT_RESULT_SETTING, UV_EXCLUDE_NEWER_SETTING, UV_INDEX_STRATEGY_SETTING, UV_PYTHON_INSTALL_MIRROR_SETTING, WORKSPACE_FAIRNESS_DURATION_SECS_SETTING, WORKSPACE_FAIRNESS_ENABLED_SETTING, WORKSPACE_FAIRNESS_MAX_PERCENT_SETTING, WORKSPACE_FAIRNESS_MIN_TOTAL_SETTING, WORKSPACE_MAX_QUEUED_JOBS_SETTING, WORKSPACE_REGISTRIES_SETTING, }, scripts::ScriptLang, stats_oss::schedule_stats, triggers::TriggerKind, utils::{ checked_worker_name, resolve_worker_suffix, Mode, GIT_VERSION, HOSTNAME, MODE_AND_ADDONS, }, worker::{ is_native_mode_from_env, reload_custom_tags_setting, validate_worker_lifecycle_env, Connection, HUB_CACHE_DIR, HUB_RT_CACHE_DIR, NATIVE_MODE_RESOLVED, TMP_LOGS_DIR, WINDMILL_DIR, WORKER_GROUP, }, KillpillSender, DEFAULT_HUB_BASE_URL, INSTANCE_NAME, METRICS_ENABLED, }; #[cfg(feature = "enterprise")] use windmill_common::worker::CLOUD_HOSTED; #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] use monitor::monitor_mem; #[cfg(any(target_os = "linux"))] use crate::cgroups::disable_oom_group; #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] use tikv_jemallocator::Jemalloc; #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] #[global_allocator] static GLOBAL: Jemalloc = Jemalloc; // Stock jemalloc only purges freed pages during alloc/free calls on app // threads, so a long-lived worker that goes quiet after a burst never runs the // purge: RSS freezes at the high-water mark and eventually OOMs under a hard // cgroup limit. Enabling the background thread makes the decay run on idle, // returning pages to the OS; the decay windows are left at jemalloc defaults // (dirty 10s, muzzy 0) on purpose — over a worker's months-long lifetime there // is no benefit to reclaiming more aggressively than that. jemalloc applies the // _RJEM_MALLOC_CONF env var after this symbol, so operators can still tune // decay or add prof:* for profiling. #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] #[allow(non_upper_case_globals)] #[export_name = "_rjem_malloc_conf"] pub static malloc_conf: &[u8] = b"background_thread:true\0"; #[cfg(feature = "parquet")] use windmill_common::global_settings::OBJECT_STORE_CONFIG_SETTING; use windmill_worker::{ get_hub_script_content_and_requirements, init_worker_internal_server_inline_utils, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, CSHARP_CACHE_DIR, DENO_CACHE_DIR, DENO_CACHE_DIR_DEPS, DENO_CACHE_DIR_NPM, GO_BIN_CACHE_DIR, GO_CACHE_DIR, JAVA_CACHE_DIR, NU_CACHE_DIR, POWERSHELL_CACHE_DIR, PY310_CACHE_DIR, PY311_CACHE_DIR, PY312_CACHE_DIR, PY313_CACHE_DIR, RUBY_CACHE_DIR, RUST_CACHE_DIR, R_CACHE_DIR, TAR_JAVA_CACHE_DIR, UV_CACHE_DIR, }; use crate::monitor::{ initial_load, load_concurrency_key_max_queued, load_disable_password_login, load_fork_workspace_tag_append_fork_suffix, load_keep_job_dir, load_metrics_debug_enabled, load_preview_tags_override, load_require_preexisting_user, load_retention_period_overrides, load_tag_per_workspace_enabled, load_tag_per_workspace_workspaces, load_workspace_fairness_duration_secs, load_workspace_fairness_enabled, load_workspace_fairness_max_percent, load_workspace_fairness_min_total, load_workspace_max_queued_jobs, monitor_db, reload_app_workspaced_route_setting, reload_audit_log_retention_days_setting, reload_base_url_setting, reload_bun_install_min_release_age_setting, reload_bunfig_install_scopes_setting, reload_critical_alert_mute_ui_setting, reload_critical_alerts_on_token_expiry_setting, reload_critical_error_channels_setting, reload_extra_pip_index_url_setting, reload_http_route_workspaced_route_setting, reload_hub_api_secret_setting, reload_hub_base_url_setting, reload_instance_events_webhook_setting, reload_job_default_timeout_setting, reload_job_isolation_setting, reload_jwt_secret_setting, reload_license_key, reload_npm_config_registry_setting, reload_nsjail_tmp_backing_setting, reload_nsjail_tmpfs_size_setting, reload_otel_tracing_proxy_setting, reload_pip_index_url_setting, reload_retention_period_setting, reload_sandbox_image_cache_max_setting, reload_sandbox_image_default_registry_setting, reload_sandbox_image_max_size_setting, reload_sandbox_image_pull_policy_setting, reload_sandbox_registry_auth_setting, reload_scim_token_setting, reload_smtp_config, reload_store_audit_logs_s3_setting, reload_uv_exclude_newer_setting, reload_uv_index_strategy_setting, reload_uv_python_install_mirror_setting, reload_worker_config, MonitorIteration, }; #[cfg(feature = "parquet")] use windmill_object_store::reload_object_store_setting; const DEFAULT_NUM_WORKERS: usize = 1; const DEFAULT_PORT: u16 = 8000; const DEFAULT_SERVER_BIND_ADDR: Ipv4Addr = Ipv4Addr::new(0, 0, 0, 0); const DEFAULT_WORKER_BIND_ADDR: Ipv4Addr = Ipv4Addr::new(127, 0, 0, 1); const BIND_ADDR_ENV: &str = "SERVER_BIND_ADDR"; #[cfg(target_os = "linux")] mod cgroups; mod db_connect; #[cfg(feature = "private")] pub mod ee; mod ee_oss; mod monitor; // Windows service support - EE feature #[cfg(all(windows, feature = "enterprise", feature = "private"))] mod windows_service_ee; pub fn setup_deno_runtime() -> anyhow::Result<()> { #[cfg(feature = "deno_core")] windmill_runtime_nativets::setup_deno_runtime()?; Ok(()) } fn update_ca_certificates_if_requested() { if std::env::var("RUN_UPDATE_CA_CERTIFICATE_AT_START") .ok() .map(|v| v.to_lowercase() == "true") .unwrap_or(false) { let ca_cert_path = std::env::var("RUN_UPDATE_CA_CERTIFICATE_PATH") .unwrap_or_else(|_| "/usr/sbin/update-ca-certificates".to_string()); println!( "RUN_UPDATE_CA_CERTIFICATE_AT_START=true, running: {}", ca_cert_path ); let output = std::process::Command::new(&ca_cert_path).output(); match output { Ok(result) => { if result.status.success() { println!("Successfully updated CA certificates"); } else { let stderr = String::from_utf8_lossy(&result.stderr); println!( "Failed to update CA certificates, but continuing startup: {}", stderr.trim() ); } } Err(e) => { println!( "Could not run update-ca-certificates command, but continuing startup: {}", e ); } } } } #[inline(always)] fn create_and_run_current_thread_inner(future: F) -> R where F: std::future::Future + 'static, R: Send + 'static, { let rt = tokio::runtime::Builder::new_multi_thread() .enable_all() .worker_threads(32) .build() .unwrap(); // Since this is the main future, we want to box it in debug mode because it tends to be fairly // large and the compiler won't optimize repeated copies. We also make this runtime factory // function #[inline(always)] to avoid holding the unboxed, unused future on the stack. #[cfg(debug_assertions)] // SAFETY: this this is guaranteed to be running on a current-thread executor let future = Box::pin(future); rt.block_on(future) } lazy_static::lazy_static! { // Period in seconds between full settings reload (12 hours by default) static ref SETTINGS_RELOAD_PERIOD_SECS: u64 = std::env::var("SETTINGS_RELOAD_PERIOD_SECS") .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(3600 * 12); // Period in seconds between polling for notify events (10s by default) static ref LISTEN_NEW_EVENTS_INTERVAL_SEC: u64 = std::env::var("LISTEN_NEW_EVENTS_INTERVAL_SEC") .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(10); } pub fn main() -> anyhow::Result<()> { // On Windows with enterprise feature, check if running as a service #[cfg(all(windows, feature = "enterprise", feature = "private"))] { if windows_service_ee::is_running_as_service() { // Run as Windows service with SCM handlers return windows_service_ee::run_as_windows_service() .map_err(|e| anyhow::anyhow!("Failed to run as Windows service: {}", e)); } } // Normal execution (console/foreground mode) setup_deno_runtime()?; create_and_run_current_thread_inner(windmill_main()) } async fn cache_hub_scripts(file_path: Option) -> anyhow::Result<()> { // The `cache` CLI mode never connects to the DB, so HUB_BASE_URL keeps its // compiled default. Allow overriding it via env so the prebuild cache step can // be pointed at a private/staging hub (e.g. a local proxy for testing). if let Ok(hub_base_url) = std::env::var("HUB_BASE_URL") { if !hub_base_url.is_empty() { tracing::info!("Overriding hub base url from env: {hub_base_url}"); windmill_common::HUB_BASE_URL.store(std::sync::Arc::new(hub_base_url)); } } let file_path = file_path.unwrap_or("./hubPaths.json".to_string()); let mut file = File::open(&file_path) .await .with_context(|| format!("Could not open {}, make sure it exists", &file_path))?; let mut contents = String::new(); file.read_to_string(&mut contents).await?; let paths = serde_json::from_str::>(&contents).with_context(|| { format!( "Could not parse {}, make sure it is a valid JSON object with string keys and values", &file_path ) })?; create_dir_all(&*HUB_CACHE_DIR)?; create_dir_all(&*BUN_BUNDLE_CACHE_DIR)?; // Ensure the backend-hardcoded git sync scripts are always cached, regardless of // hubPaths.json contents. These are run by backend-driven sync (not necessarily listed // in hubPaths.json), so airgapped workers would otherwise miss them on a cache lookup. let mut all_paths: Vec = paths.into_values().collect(); for git_sync_path in [ windmill_common::workspaces::LATEST_GIT_SYNC_SCRIPT_PATH, windmill_common::workspaces::GIT_SYNC_PULL_SCRIPT_PATH, ] { let git_sync_path = git_sync_path.to_string(); if !all_paths.contains(&git_sync_path) { all_paths.push(git_sync_path); } } for path in &all_paths { tracing::info!("Caching hub script at {path}"); let res = get_hub_script_content_and_requirements(Some(path), None).await?; if res .language .as_ref() .is_some_and(|x| x == &ScriptLang::Deno) { let job_dir = format!("{}/cache_init/{}", *WINDMILL_DIR, Uuid::new_v4()); create_dir_all(&job_dir)?; let _ = windmill_worker::generate_deno_lock( &Uuid::nil(), &res.content, &mut 0, &mut None, &job_dir, None, "global", "global", "", &mut None, ) .await?; tokio::fs::remove_dir_all(job_dir).await?; } else if res.language.as_ref().is_some_and(|x| x == &ScriptLang::Bun) { let job_id = Uuid::new_v4(); let job_dir = format!("{}/cache_init/{}", *WINDMILL_DIR, job_id); create_dir_all(&job_dir)?; if let Some(lock) = res.lockfile { // The hub occasionally returns a malformed `lockfile` field — e.g. the // raw script source instead of the expected `\n//bun.lockb\n` // shape. A valid lockfile always starts with the package.json (`{...}`), // so anything else is bogus and would make `bun install` choke trying to // parse TypeScript as JSON. Skip those rather than aborting the entire cache. if !lock.trim_start().starts_with('{') { tracing::warn!( "Hub script {path} returned a malformed lockfile (does not start with a package.json object), skipping prebundling" ); } else { let _ = windmill_worker::prepare_job_dir(&lock, &job_dir).await?; let envs = windmill_worker::get_common_bun_proc_envs(None).await; if let Err(e) = windmill_worker::install_bun_lockfile( &mut 0, &mut None, &job_id, "admins", None, &job_dir, "cache_init", envs.clone(), false, &mut None, false, ) .await { // A single broken hub script (malformed lockfile, missing dep, …) // shouldn't abort the entire cache run — log and move on. tracing::error!( "Failed to install lockfile for hub script {path}, skipping: {e:#}" ); } else if let Err(e) = windmill_worker::prebundle_bun_script( &res.content, &lock, &path, &job_id, "admins", None, &job_dir, "", "cache_init", "", &mut None, &None, ) .await { tracing::error!("Failed to prebundle hub script {path}, skipping: {e:#}"); } } } else { tracing::warn!("No lockfile found for bun script {path}, skipping..."); } tokio::fs::remove_dir_all(job_dir).await?; } } Ok(()) } /// Raw resource type from hub API (schema is a JSON string) #[derive(serde::Deserialize)] struct HubResourceTypeRaw { pub id: i64, pub name: String, pub schema: Option, pub app: String, pub description: Option, } /// Processed resource type with parsed schema #[derive(serde::Deserialize, serde::Serialize, Clone)] pub struct HubResourceType { pub id: i64, pub name: String, pub schema: Option, pub app: String, pub description: Option, } const HUB_RT_CACHE_FILE: &str = "resource_types.json"; async fn cache_hub_resource_types() -> anyhow::Result<()> { println!("Caching resource types from hub..."); let response = HTTP_CLIENT .get(format!("{}/resource_types/list", DEFAULT_HUB_BASE_URL)) .header("Accept", "application/json") .send() .await .with_context(|| "Failed to fetch resource types from hub")?; if !response.status().is_success() { anyhow::bail!( "Failed to fetch resource types from hub: {}", response.status() ); } let raw_types: Vec = response .json::>() .await .with_context(|| "Failed to parse resource types from hub")?; // Parse schema strings into JSON values let resource_types: Vec = raw_types .into_iter() .filter_map(|rt| { let schema = match rt.schema { Some(s) => match serde_json::from_str(&s) { Ok(v) => Some(v), Err(e) => { println!("Warning: failed to parse schema for {}: {}", rt.name, e); return None; } }, None => None, }; Some(HubResourceType { id: rt.id, name: rt.name, schema, app: rt.app, description: rt.description, }) }) .collect(); println!("Fetched {} resource types from hub", resource_types.len()); create_dir_all(&*HUB_RT_CACHE_DIR)?; let cache_path = format!("{}/{}", *HUB_RT_CACHE_DIR, HUB_RT_CACHE_FILE); let content = serde_json::to_string_pretty(&resource_types) .with_context(|| "Failed to serialize resource types")?; std::fs::write(&cache_path, content) .with_context(|| format!("Failed to write cache file to {}", cache_path))?; println!("Cached resource types to {}", cache_path); Ok(()) } pub async fn sync_cached_resource_types(db: &sqlx::Pool) -> anyhow::Result<()> { let cache_path = format!("{}/{}", *HUB_RT_CACHE_DIR, HUB_RT_CACHE_FILE); if tokio::fs::metadata(&cache_path).await.is_err() { tracing::info!( "No cached resource types found at {}, skipping sync", cache_path ); return Ok(()); } tracing::info!("Syncing cached resource types to admins workspace..."); let content = tokio::fs::read_to_string(&cache_path) .await .with_context(|| format!("Failed to read cache file from {}", cache_path))?; let cached_types: Vec = serde_json::from_str(&content).with_context(|| "Failed to parse cached resource types")?; tracing::info!("Found {} cached resource types", cached_types.len()); // Get existing resource types in admins workspace let existing_types: Vec<(String, Option, Option)> = sqlx::query_as( "SELECT name, schema, description FROM resource_type WHERE workspace_id = 'admins'", ) .fetch_all(db) .await .with_context(|| "Failed to fetch existing resource types")?; let existing_map: std::collections::HashMap< String, (Option, Option), > = existing_types .into_iter() .map(|(name, schema, desc)| (name, (schema, desc))) .collect(); let mut synced_count = 0; let mut skipped_count = 0; for rt in cached_types { // Check if resource type already exists with same schema and description if let Some((existing_schema, existing_desc)) = existing_map.get(&rt.name) { if existing_schema == &rt.schema && existing_desc == &rt.description { skipped_count += 1; continue; } } // Insert or update resource type sqlx::query( "INSERT INTO resource_type (workspace_id, name, schema, description, edited_at) VALUES ('admins', $1, $2, $3, now()) ON CONFLICT (workspace_id, name) DO UPDATE SET schema = EXCLUDED.schema, description = EXCLUDED.description, edited_at = now()", ) .bind(&rt.name) .bind(&rt.schema) .bind(&rt.description) .execute(db) .await .with_context(|| format!("Failed to upsert resource type {}", rt.name))?; synced_count += 1; } tracing::info!( "Synced {} resource types to admins workspace ({} skipped as unchanged)", synced_count, skipped_count ); Ok(()) } fn print_help() { println!("Windmill - a fast, open-source workflow engine and job runner."); println!(); println!("Usage:"); println!(" windmill [SUBCOMMAND]"); println!(); println!("Subcommands:"); println!(" help | -h | --help Show this help information and exit"); println!(" version Show Windmill version and exit"); println!(" cache [hubPaths.json] Pre-cache hub scripts (default: ./hubPaths.json)"); println!(" cache-rt Pre-cache hub resource types"); println!(" sync-config Sync instance config from a YAML file to the database"); println!(" operator Run the Kubernetes operator (watches a ConfigMap)"); println!(); println!("Environment variables (name = default):"); println!(" DATABASE_URL = The Postgres database url."); println!(" DATABASE_URL_FILE = None Read the database url from a file instead, e.g. a mounted secret (takes precedence over DATABASE_URL)"); println!(" MODE = standalone Mode: standalone | worker | server | agent"); println!(" BASE_URL = http://localhost:8000 Public base URL of your instance (overridden by instance settings)"); println!( " PORT = {} HTTP port (server/indexer/MCP modes)", DEFAULT_PORT ); println!( " SERVER_BIND_ADDR = IP to bind to (server: {}, worker: {})", DEFAULT_SERVER_BIND_ADDR, DEFAULT_WORKER_BIND_ADDR ); println!( " NUM_WORKERS = {} Number of workers (standalone/worker modes)", DEFAULT_NUM_WORKERS ); println!(" WORKER_GROUP = default Worker group this worker belongs to",); println!(" JSON_FMT = false Output logs in JSON instead of logfmt"); println!(" METRICS_ADDR = None (EE only) Prometheus metrics addr at /metrics; set \"true\" to use :8001"); println!(" SUPERADMIN_SECRET = None Virtual superadmin token (server)"); println!(" NO_AUTH = false Bypass all auth; every request acts as the admin@windmill.dev superadmin (only behind a trusted gateway; ignored when CLOUD_HOSTED)"); println!(" LICENSE_KEY = None (EE only) Enterprise license key (workers require valid key)"); println!(" RUN_UPDATE_CA_CERTIFICATE_AT_START = false Run system CA update at startup"); println!(" RUN_UPDATE_CA_CERTIFICATE_PATH = /usr/sbin/update-ca-certificates Path to CA update tool"); println!(" SYNC_CACHED_RT = false Sync cached resource types to admins workspace on server start"); println!(" HUB_BASE_URL = https://hub.windmill.dev Hub to fetch scripts from in `cache` mode (server/worker use the DB setting instead)"); println!(); println!("Notes:"); println!("- Advanced and less commonly used settings are managed via the database and are omitted here."); println!("- At startup, Windmill logs currently set configuration keys for visibility."); } async fn windmill_main() -> anyhow::Result<()> { let (killpill_tx, mut killpill_rx) = KillpillSender::new(2); let mut monitor_killpill_rx = killpill_tx.subscribe(); let (killpill_phase2_tx, _killpill_phase2_rx) = tokio::sync::broadcast::channel::<()>(2); let server_killpill_rx = killpill_phase2_tx.subscribe(); let shutdown_tx = killpill_tx.clone(); let shutdown_rx = killpill_tx.subscribe(); tokio::spawn(async move { if let Err(e) = windmill_common::shutdown_signal(shutdown_tx, shutdown_rx).await { tracing::error!("Error in shutdown signal: {e:#}"); } }); dotenv::dotenv().ok(); update_ca_certificates_if_requested(); if std::env::var("RUST_LOG").is_err() { unsafe { std::env::set_var("RUST_LOG", "info") } } if let Err(_e) = rustls::crypto::ring::default_provider().install_default() { println!("Failed to install rustls crypto provider"); } #[cfg(feature = "enterprise")] if *CLOUD_HOSTED { // Block access to AWS/GCP metadata endpoints for security in cloud-hosted mode. // This is a best-effort attempt; if it fails, just warn and continue. if let Err(e) = std::process::Command::new("sh") .arg("-c") .arg("iptables -A OUTPUT -d 169.254.169.254 -j DROP && iptables -A FORWARD -d 169.254.169.254 -j DROP") .status() { println!("Failed to run iptables to block metadata endpoint: {e}"); } else { println!("Successfully blocked metadata endpoint using iptables"); } } let hostname = HOSTNAME.to_owned(); let mode_and_addons = MODE_AND_ADDONS.clone(); let mode = mode_and_addons.mode; if mode == Mode::Standalone { println!("Running in standalone mode"); } else if mode == Mode::MCP { println!("Running in MCP mode"); } if *windmill_common::worker::NO_AUTH { println!("############################################################"); println!("# NO_AUTH mode is ENABLED: authentication is fully #"); println!("# bypassed and every request is treated as the #"); println!("# admin@windmill.dev superadmin. Only run this behind a #"); println!("# trusted authenticating gateway on a private network. #"); println!("############################################################"); } #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] println!("jemalloc enabled"); let cli_arg = std::env::args().nth(1).unwrap_or_default(); match cli_arg.as_str() { "-h" | "--help" | "help" => { print_help(); return Ok(()); } "cache" => { tracing_subscriber::fmt::init(); #[cfg(feature = "embedding")] { println!("Caching embedding model..."); windmill_api::embeddings::ModelInstance::load_model_files().await?; println!("Cached embedding model"); } #[cfg(not(feature = "embedding"))] { println!("Embeddings are not enabled, ignoring..."); } cache_hub_scripts(std::env::args().nth(2)).await?; return Ok(()); } "-v" | "--version" | "version" => { println!("Windmill {}", GIT_VERSION); return Ok(()); } "prepare-deps" => { // CLI command for preparing dependencies without database access // Used by the debugger to install dependencies for scripts windmill_worker::run_prepare_deps_cli().await?; return Ok(()); } "cache-rt" => { cache_hub_resource_types().await?; return Ok(()); } "sync-config" => { tracing_subscriber::fmt::init(); let path = std::env::args().nth(2).unwrap_or_else(|| { eprintln!("Usage: windmill sync-config "); std::process::exit(1); }); let contents = tokio::fs::read_to_string(&path) .await .with_context(|| format!("Could not read config file: {path}"))?; let mut config: windmill_common::instance_config::InstanceConfig = serde_yml::from_str(&contents) .with_context(|| format!("Could not parse YAML from: {path}"))?; windmill_common::instance_config::resolve_env_refs(&mut config.global_settings) .map_err(|var| anyhow::anyhow!("environment variable '{var}' not found"))?; tracing::info!("Connecting to database..."); let db = crate::db_connect::initial_connection().await?; config.sync_to_db(&db).await?; tracing::info!("Synced instance config from {path}"); return Ok(()); } #[cfg(feature = "operator")] "operator" => { tracing_subscriber::fmt::init(); tracing::info!("Starting Windmill Kubernetes operator..."); tracing::info!("Connecting to database..."); #[cfg(all(feature = "enterprise", feature = "private"))] let (operator_killpill_tx, operator_killpill_rx) = tokio::sync::broadcast::channel::<()>(2); let db = crate::db_connect::operator_connection( #[cfg(all(feature = "enterprise", feature = "private"))] operator_killpill_rx, ) .await?; #[cfg(all(feature = "enterprise", feature = "private"))] tokio::spawn(async move { if let Ok(()) = tokio::signal::ctrl_c().await { let _ = operator_killpill_tx.send(()); } }); tracing::info!("Database connected. Starting ConfigMap watcher..."); windmill_operator::run(db).await?; return Ok(()); } _ => {} } #[allow(unused_mut)] let mut num_workers = if mode == Mode::Server || mode == Mode::Indexer || mode == Mode::MCP { 0 } else if is_native_mode_from_env() { NATIVE_MODE_RESOLVED.store(true, std::sync::atomic::Ordering::Relaxed); println!("Native mode enabled: forcing NUM_WORKERS=8"); 8 } else { std::env::var("NUM_WORKERS") .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(DEFAULT_NUM_WORKERS as i32) }; if num_workers > 1 && !is_native_mode_from_env() { if std::env::var("I_ACK_NUM_WORKERS_IS_UNSAFE").is_ok_and(|x| x == "1" || x == "true") { println!( "WARNING: Running with NUM_WORKERS={} without native mode. \ This is not recommended. Use at your own risk.", num_workers ); } else { eprintln!( "WARNING: NUM_WORKERS={} > 1 is only safe for native workers. \ Falling back to NUM_WORKERS=1. Set NATIVE_MODE=true for native-only workers.", num_workers ); num_workers = 1; } } validate_worker_lifecycle_env()?; let server_mode = !std::env::var("DISABLE_SERVER") .ok() .and_then(|x| x.parse::().ok()) .unwrap_or(false) && (mode == Mode::Server || mode == Mode::Standalone); let indexer_mode = mode == Mode::Indexer; let mcp_mode = mode == Mode::MCP; let default_bind_addr = if server_mode || indexer_mode || mcp_mode { DEFAULT_SERVER_BIND_ADDR } else { DEFAULT_WORKER_BIND_ADDR }; let server_bind_address: IpAddr = std::env::var(BIND_ADDR_ENV) .ok() .and_then(|x| x.parse().ok()) .unwrap_or(IpAddr::from(default_bind_addr)); let (conn, first_suffix, agent_config) = if mode == Mode::Agent { let agent_config = match AgentConfig::from_env() { Ok(config) => config, Err(e) => { tracing::error!("{e}"); std::process::exit(1); } }; tracing::info!( "Creating http client for cluster using base internal url {}", agent_config.base_internal_url ); let suffix = resolve_worker_suffix(&hostname, 1)?; ( Connection::Http(agent_config.build_http_client(&suffix)), Some(suffix), Some(agent_config), ) } else { println!("Connecting to database..."); let db = crate::db_connect::initial_connection().await?; let num_version = sqlx::query_scalar!("SELECT version()").fetch_one(&db).await; println!( "PostgreSQL version: {} (windmill require PG >= 14)", num_version .ok() .flatten() .unwrap_or_else(|| "UNKNOWN".to_string()) ); // Load OTEL tracing proxy settings and initialize deno_telemetry if nativets tracing is enabled // This must happen before any Deno runtime is created #[cfg(all(feature = "private", feature = "enterprise"))] { reload_otel_tracing_proxy_setting(&Connection::Sql(db.clone())).await; #[cfg(feature = "deno_core")] if windmill_worker::is_otel_tracing_proxy_enabled_for_lang(&ScriptLang::Nativets).await { match windmill_worker::load_internal_otel_exporter().await { Ok(()) => { tracing::info!("Internal OTEL exporter initialized for nativets tracing"); } Err(e) => { tracing::error!("Failed to initialize internal OTEL exporter: {}", e); } } } } load_otel(&db).await; println!("Database connected"); (Connection::Sql(db), None, None) }; let environment = if let Ok(environment) = std::env::var("OTEL_ENVIRONMENT") { environment } else { load_base_url(&conn) .await .unwrap_or_else(|_| "local".to_string()) .trim_start_matches("https://") .trim_start_matches("http://") .split(".") .next() .unwrap_or_else(|| "local") .to_string() }; let _guard = windmill_common::tracing_init::initialize_tracing(&hostname, &mode, &environment); let is_agent = mode == Mode::Agent; let mut migration_handle: Option> = None; #[cfg(feature = "parquet")] let disable_s3_store = std::env::var("DISABLE_S3_STORE") .ok() .is_some_and(|x| x == "1" || x == "true"); if let Some(db) = conn.as_sql() { if !is_agent && !indexer_mode && !mcp_mode { let skip_migration = std::env::var("SKIP_MIGRATION") .map(|val| val == "true") .unwrap_or(false); if !skip_migration { if mode == Mode::Worker { windmill_api::wait_for_db_migrations(&db, killpill_rx.resubscribe()).await?; } else { migration_handle = windmill_api::migrate_db(&db, killpill_rx.resubscribe()).await?; } } else { tracing::info!("SKIP_MIGRATION set, skipping db migration...") } // Sync cached resource types to admins workspace if SYNC_CACHED_RT is set if std::env::var("SYNC_CACHED_RT") .ok() .map(|v| v.to_lowercase() == "true" || v == "1") .unwrap_or(false) { if let Err(e) = sync_cached_resource_types(db).await { tracing::warn!("Failed to sync cached resource types: {:#}", e); } } } } if killpill_rx.try_recv().is_ok() { tracing::info!("Received early killpill, aborting startup"); return Ok(()); } let worker_mode = num_workers > 0; if worker_mode { #[cfg(any(target_os = "linux"))] if let Err(e) = disable_oom_group() { tracing::warn!( "Failed to disable cgroup OOM group kill: {e:?}. \ When a job exceeds memory, the OOM killer will kill the entire pod \ instead of just the offending job process" ); } // Lower the worker's oom_score_adj so the OOM killer strongly prefers killing // job subprocesses (oom_score_adj=JOB_OOM_SCORE_ADJ) over the worker itself. // Kubernetes sets it high for burstable QoS (e.g. 937), leaving a tiny gap vs jobs. // Requires CAP_SYS_RESOURCE to lower it; if missing, we just warn. #[cfg(any(target_os = "linux"))] { // Badness is (memory used, in permille of host RAM) + oom_score_adj, so the gap // must exceed the worker's own footprint in permille to actually steer the kill. // 100 covers a worker holding up to ~10% of host RAM. const MIN_OOM_SCORE_GAP: i32 = 100; let job_adj = *windmill_common::worker::JOB_OOM_SCORE_ADJ; match std::fs::read_to_string("/proc/self/oom_score_adj") { Ok(current) => { let current = current.trim().to_string(); match current.parse::() { Ok(mut worker_adj) => { if worker_adj > 0 { match std::fs::write("/proc/self/oom_score_adj", "0") { Ok(_) => { tracing::info!( "Lowered worker oom_score_adj from {worker_adj} to 0" ); worker_adj = 0; } Err(e) => { tracing::warn!( "Could not lower worker oom_score_adj from {worker_adj} to 0: {e}. \ Add CAP_SYS_RESOURCE to the container to fix this" ); } } } let gap = job_adj - worker_adj; if gap >= MIN_OOM_SCORE_GAP { tracing::info!( "Worker oom_score_adj={worker_adj}, jobs get {job_adj} (gap={gap})" ); } else { tracing::warn!( "Worker oom_score_adj={worker_adj}, jobs get {job_adj} (gap={gap}): \ too small to reliably steer the OOM killer to the job. \ Raise JOB_OOM_SCORE_ADJ or lower the worker's own score" ); } } Err(e) => { tracing::warn!( "Could not parse worker oom_score_adj '{current}': {e}. \ Cannot tell whether jobs (oom_score_adj={job_adj}) outrank the worker" ); } } } Err(e) => { tracing::warn!("Could not read worker oom_score_adj: {e}"); } } } } // Resolve native mode early (before connect_db) so connection pool size accounts for it. // native_mode can come from env OR from the DB worker group config. if worker_mode && !is_native_mode_from_env() { if let Some(db) = conn.as_sql() { let native_from_db: bool = sqlx::query_scalar!( "SELECT (config->>'native_mode')::boolean FROM config WHERE name = $1", format!("worker__{}", *windmill_common::worker::WORKER_GROUP) ) .fetch_optional(db) .await .ok() .flatten() .flatten() .unwrap_or(false); if native_from_db { NATIVE_MODE_RESOLVED.store(true, std::sync::atomic::Ordering::Relaxed); num_workers = 8; tracing::info!( "Native mode detected from worker config (early): forcing NUM_WORKERS=8" ); } } } let conn = if mode == Mode::Agent { conn } else { // Drop the initial connection pool before creating the main one. // With low PostgreSQL max_connections, both pools existing simultaneously // can exhaust all available connection slots, causing connect_db to hang. drop(conn); let db = crate::db_connect::connect_db( server_mode, indexer_mode, worker_mode, num_workers, #[cfg(feature = "private")] killpill_rx.resubscribe(), ) .await?; // NOTE: Variable/resource cache initialization moved to API server in windmill-api Connection::Sql(db) }; #[cfg(feature = "enterprise")] tracing::info!( " ############################## Windmill Enterprise Edition {GIT_VERSION} ##############################" ); #[cfg(not(feature = "enterprise"))] tracing::info!( " ############################## Windmill Community Edition {GIT_VERSION} ##############################" ); display_config(&ENV_SETTINGS); #[cfg(feature = "enterprise")] { // load the license key and check if it's valid // if not valid and not server mode just quit // if not expired and server mode then force renewal // if key still invalid and num_workers > 0, set to 0 if let Err(err) = reload_license_key(&conn).await { tracing::error!("Failed to reload license key: {err:#}"); if is_agent { tracing::error!( "Agent worker cannot connect to server. Please check AGENT_TOKEN and BASE_INTERNAL_URL" ); std::process::exit(1); } } let valid_key = LICENSE_KEY_VALID.load(std::sync::atomic::Ordering::Relaxed); if !valid_key && !server_mode { tracing::error!("Invalid license key, workers require a valid license key"); } if server_mode || mcp_mode { if let Some(db) = conn.as_sql() { // only force renewal if invalid but not empty (= expired) let renewed_now = maybe_renew_license_key_on_start( &HTTP_CLIENT, &db, !valid_key && !LICENSE_KEY_ID.load().is_empty(), ) .await; if renewed_now { if let Err(err) = reload_license_key(&conn).await { tracing::error!("Failed to reload license key: {err:#}"); } } } else { panic!("Server mode requires a database connection"); } } } if server_mode || worker_mode || indexer_mode || mcp_mode { let port_var = std::env::var("PORT") .or_else(|_| std::env::var("BACKEND_PORT")) .ok() .and_then(|x| x.parse().ok()); let port = if server_mode || indexer_mode || mcp_mode { port_var.unwrap_or(DEFAULT_PORT as u16) } else { port_var.unwrap_or(0) }; let default_base_internal_url = format!("http://localhost:{}", port.to_string()); // since it's only on server mode, the port is statically defined let base_internal_url: String = if let Ok(base_url) = std::env::var("BASE_INTERNAL_URL") { if !is_agent { tracing::warn!("BASE_INTERNAL_URL is now unecessary and ignored unless the mode is 'agent', you can remove it."); default_base_internal_url.clone() } else { base_url } } else { default_base_internal_url.clone() }; initial_load( &conn, killpill_tx.clone(), worker_mode, server_mode, #[cfg(feature = "parquet")] disable_s3_store, ) .await; monitor_db( &conn, &base_internal_url, server_mode, worker_mode, true, killpill_tx.clone(), None, ) .await; #[cfg(feature = "prometheus")] if let Some(db) = conn.as_sql() { crate::monitor::monitor_pool(&db).await; } if let Some(db) = conn.as_sql() { crate::monitor::monitor_pool_otel(&db).await; } send_logs_to_object_store(&conn, &hostname, &mode); #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] if !worker_mode { monitor_mem().await; } let addr = SocketAddr::from((server_bind_address, port)); let listener = tokio::net::TcpListener::bind(addr) .await .context("binding main windmill server")?; let (base_internal_tx, base_internal_rx) = tokio::sync::oneshot::channel::(); DirBuilder::new() .recursive(true) .create(&*WINDMILL_DIR) .expect("could not create initial server dir"); #[cfg(feature = "tantivy")] let should_index_jobs = mode == Mode::Indexer || mode_and_addons.indexer; #[cfg(feature = "tantivy")] if should_index_jobs { if let Some(db) = conn.as_sql() { reload_indexer_config(&db).await; } } #[cfg(feature = "tantivy")] let (index_reader, index_writer) = if should_index_jobs { if let Some(db) = conn.as_sql() { let mut indexer_rx = killpill_rx.resubscribe(); let (mut reader, mut writer) = (None, None); tokio::select! { _ = indexer_rx.recv() => { tracing::info!("Received killpill, aborting index initialization"); }, res = windmill_indexer::completed_runs_oss::init_index(&db) => { let res = res?; if let Some(r) = res { reader = Some(r.0); writer = Some(r.1); } } } (reader, writer) } else { (None, None) } } else { (None, None) }; #[cfg(feature = "tantivy")] let indexer_f = { let indexer_rx = killpill_rx.resubscribe(); let index_writer2 = index_writer.clone(); async { if let Some(db) = conn.as_sql() { if let Some(index_writer) = index_writer2 { windmill_indexer::completed_runs_oss::run_indexer( db.clone(), index_writer, indexer_rx, ) .await?; } } Ok(()) } }; #[cfg(all(feature = "tantivy", feature = "parquet"))] let (log_index_reader, log_index_writer) = if should_index_jobs { if let Some(db) = conn.as_sql() { let mut indexer_rx = killpill_rx.resubscribe(); let (mut reader, mut writer) = (None, None); tokio::select! { _ = indexer_rx.recv() => { tracing::info!("Received killpill, aborting index initialization"); }, res = windmill_indexer::service_logs_oss::init_index(&db, killpill_tx.clone()) => { let res = res?; if let Some(r) = res { reader = Some(r.0); writer = Some(r.1); } } } (reader, writer) } else { (None, None) } } else { (None, None) }; #[cfg(all(feature = "tantivy", feature = "parquet"))] let log_indexer_f = { let log_indexer_rx = killpill_rx.resubscribe(); let log_index_writer2 = log_index_writer.clone(); async { if let Some(db) = conn.as_sql() { if let Some(log_index_writer) = log_index_writer2 { windmill_indexer::service_logs_oss::run_indexer( db.clone(), log_index_writer, log_indexer_rx, ) .await?; } } Ok(()) } }; #[cfg(not(feature = "tantivy"))] let index_reader = None; #[cfg(not(feature = "tantivy"))] let indexer_f = async { Ok(()) as anyhow::Result<()> }; #[cfg(not(all(feature = "tantivy", feature = "parquet")))] let log_index_reader = None; #[cfg(not(all(feature = "tantivy", feature = "parquet")))] let log_indexer_f = async { Ok(()) as anyhow::Result<()> }; // Resubscribe for OTEL tracing proxy before workers_f captures killpill_rx #[cfg(all(feature = "private", feature = "enterprise"))] let otel_killpill_rx = killpill_rx.resubscribe(); let server_f = async { if !is_agent { if let Some(db) = conn.as_sql() { windmill_api::run_server( db.clone(), index_reader, log_index_reader, listener, server_killpill_rx, base_internal_tx, server_mode, mcp_mode, base_internal_url.clone(), None, ) .await?; } } else { base_internal_tx .send(base_internal_url.clone()) .map_err(|e| { anyhow::anyhow!("Could not send base_internal_url to agent: {e:#}") })?; } Ok(()) as anyhow::Result<()> }; let workers_f = async { let mut rx = killpill_rx.resubscribe(); if !killpill_rx.try_recv().is_ok() { let base_internal_url = base_internal_rx.await?; if worker_mode { let worker_internal_server_killpill_rx = killpill_rx.resubscribe(); init_worker_internal_server_inline_utils( worker_internal_server_killpill_rx, base_internal_url.clone(), )?; let mut workers = vec![]; for i in 0..num_workers { let suffix = if i == 0 && first_suffix.is_some() { first_suffix.as_ref().unwrap().clone() } else { resolve_worker_suffix(&hostname, i as usize + 1)? }; let worker_conn = WorkerConn { conn: if i == 0 || mode != Mode::Agent { conn.clone() } else { Connection::Http( agent_config .as_ref() .expect("agent_config must be set in agent mode") .build_http_client(&suffix), ) }, worker_name: checked_worker_name( mode == Mode::Agent, WORKER_GROUP.as_str(), &suffix, )?, }; workers.push(worker_conn); } run_workers( rx, killpill_tx.clone(), base_internal_url.clone(), hostname.clone(), &workers, ) .await?; tracing::info!("All workers exited."); killpill_tx.send(); } else { rx.recv().await?; } } if killpill_phase2_tx.receiver_count() > 0 { if worker_mode { tracing::info!("Starting phase 2 of shutdown"); } killpill_phase2_tx.send(())?; if worker_mode { tracing::info!("Phase 2 of shutdown completed"); } } Ok(()) }; let monitor_f = async { let tx = killpill_tx.clone(); let conn = conn.clone(); match conn { Connection::Sql(ref db) => { let base_internal_url = base_internal_url.to_string(); let db = db.clone(); let h = tokio::spawn(async move { // Initialize last_event_id to current max to avoid processing old events on startup let mut last_event_id: i64 = match windmill_common::notify_events::get_latest_event_id(&db).await { Ok(id) => { tracing::info!( "Initialized notify event polling with last_event_id: {}", id ); id } Err(e) => { tracing::warn!( "Could not get latest event id, starting from 0: {e:#}" ); 0 } }; let mut last_settings_reload = Instant::now(); let mut monitor_iteration: u64 = 0; let rd_shift: u8 = rand::rng().random_range(0..200); loop { let db = db.clone(); tokio::select! { biased; Some(_) = async { if let Some(jh) = migration_handle.take() { tracing::info!("migration job finished"); Some(jh.await) } else { None }} => { continue; }, _ = monitor_killpill_rx.recv() => { tracing::info!("received killpill for monitor job"); break; }, _ = tokio::time::sleep(Duration::from_secs(*LISTEN_NEW_EVENTS_INTERVAL_SEC)) => { // Poll for new events from notify_event table match windmill_common::notify_events::poll_notify_events(&db, last_event_id).await { Ok(events) => { for event in events { if !*windmill_common::QUIET_LOGS { tracing::info!("Processing notify event: channel={}, payload={}", event.channel, event.payload); } process_notify_event( &event.channel, &event.payload, &db, &conn, &tx, server_mode, worker_mode, #[cfg(feature = "parquet")] disable_s3_store, ).await; last_event_id = last_event_id.max(event.id); } } Err(e) => { tracing::error!("Error polling notify events: {e:#}"); } } // Periodic full settings reload if last_settings_reload.elapsed() > Duration::from_secs(*SETTINGS_RELOAD_PERIOD_SECS) { tracing::info!("Reloading settings and license key after {}s", Duration::from_secs(*SETTINGS_RELOAD_PERIOD_SECS).as_secs()); initial_load( &conn, tx.clone(), worker_mode, server_mode, #[cfg(feature = "parquet")] disable_s3_store, ) .await; #[cfg(feature = "enterprise")] if let Err(err) = reload_license_key(&conn).await { tracing::error!("Failed to reload license key: {err:#}"); } last_settings_reload = Instant::now(); } let monitor_start = Instant::now(); let warn_handle = if server_mode { Some(tokio::spawn(async move { tokio::time::sleep(Duration::from_secs(5)).await; tracing::warn!("monitor task has been running for more than 5s"); })) } else { None }; // Hard cap on a single monitor pass. monitor_db runs all its // periodic tasks under one join!, so a single task stuck on a // non-DB await (statement_timeout only bounds DB statements) // would otherwise freeze the whole loop indefinitely — silently // stopping critical maintenance like audit-partition creation. // Larger than statement_timeout (5min) so a slow-but-progressing // statement is never killed prematurely. const MONITOR_DB_TIMEOUT: Duration = Duration::from_secs(600); let monitor_timed_out = tokio::time::timeout( MONITOR_DB_TIMEOUT, monitor_db( &conn, &base_internal_url, server_mode, worker_mode, false, tx.clone(), Some(MonitorIteration { rd_shift, iter: monitor_iteration, }), ), ) .await .is_err(); if monitor_timed_out { windmill_common::utils::report_critical_error( format!( "monitor task did not finish within {}s and was aborted; \ a background maintenance task is likely stuck. \ Continuing to the next iteration.", MONITOR_DB_TIMEOUT.as_secs() ), db.clone(), None, None, ) .await; } monitor_iteration += 1; if let Some(handle) = warn_handle { handle.abort(); } let elapsed = monitor_start.elapsed(); if server_mode && elapsed >= Duration::from_secs(5) { tracing::info!("monitor task finished in {elapsed:.1?}"); } }, } } }); if let Err(e) = h.await { tracing::error!("Error waiting for monitor handle: {e:#}") } } ref conn @ Connection::Http(_) => { pub const RELOAD_FREQUENCY: Duration = Duration::from_secs(12 * 60 * 60); let mut last_time_config_reload: Instant = Instant::now(); loop { tokio::select! { _ = monitor_killpill_rx.recv() => { tracing::info!("Received killpill, exiting"); break; }, _ = tokio::time::sleep(Duration::from_secs(30)) => { // Reload config every 12h if last_time_config_reload.elapsed() > RELOAD_FREQUENCY { last_time_config_reload = Instant::now(); tracing::info!("Reloading config after 12 hours"); initial_load(&conn, tx.clone(), worker_mode, server_mode, #[cfg(feature = "parquet")] disable_s3_store).await; if let Err(e) = reload_license_key(&conn).await { tracing::error!("Failed to reload license key on agent: {e:#}"); } #[cfg(feature = "enterprise")] ee_oss::verify_license_key(conn.as_sql()).await; } // self-throttled; picks up a key fixed on the server without // waiting for the 12h reload #[cfg(feature = "enterprise")] crate::monitor::refetch_license_key_if_invalid(conn).await; // update min version explicitly. // for sql connection it is the part of monitor_db. // TODO: pass worker names for min keep-alive alerts (for HTTP connection) windmill_common::min_version::update_min_version(conn, true, vec![], false).await; } }; } } }; tracing::info!("Monitor exited"); killpill_tx.send(); Ok(()) as anyhow::Result<()> }; let metrics_f = async { let enabled = METRICS_ENABLED.load(std::sync::atomic::Ordering::Relaxed); #[cfg(not(all(feature = "enterprise", feature = "prometheus")))] if enabled { tracing::error!("Metrics are only available in the EE, ignoring..."); } #[cfg(all(feature = "enterprise", feature = "prometheus"))] if let Err(e) = windmill_common::serve_metrics( *windmill_common::METRICS_ADDR, _killpill_phase2_rx, num_workers > 0, enabled, ) .await { tracing::error!("Error serving metrics: {e:#}"); } Ok(()) as anyhow::Result<()> }; let otel_tracing_proxy_f = async { #[cfg(all(feature = "private", feature = "enterprise"))] if worker_mode { if let Some(db) = conn.as_sql() { if let Err(e) = windmill_worker::start_jobs_otel_tracing( db.clone(), otel_killpill_rx, num_workers, ) .await { tracing::error!("Jobs OTEL tracing error: {}", e); } } } Ok(()) as anyhow::Result<()> }; if server_mode { if let Some(db) = conn.as_sql() { schedule_stats(&db, &HTTP_CLIENT).await; } } // `workers_f` must stay ahead of `server_f`: these are polled on one task in // declaration order, and `run_server` yields once after handing over the base // internal url so the workers get past that oneshot before it builds its router. // Ordering `server_f` first makes them wait out the whole build instead. if mcp_mode { futures::try_join!(workers_f, server_f)?; } else { futures::try_join!( workers_f, monitor_f, server_f, metrics_f, otel_tracing_proxy_f, indexer_f, log_indexer_f )?; } } else { tracing::info!("Nothing to do, exiting."); } send_current_log_file_to_object_store(&conn, &hostname, &mode).await; if let Some(db) = conn.as_sql() { tracing::info!("Exiting connection pool"); tokio::select! { _ = db.close() => { tracing::info!("Database connection pool closed"); }, _ = tokio::time::sleep(Duration::from_secs(15)) => { tracing::warn!("Could not close database connection pool in time (15s). Exiting anyway."); } } } std::process::exit(0); } /// Process a single notify event from the polling-based event system. /// This replaces the old PgListener notification handling. #[allow(unused_variables)] async fn process_notify_event( channel: &str, payload: &str, db: &Pool, conn: &Connection, tx: &KillpillSender, server_mode: bool, worker_mode: bool, #[cfg(feature = "parquet")] disable_s3_store: bool, ) { match channel { "notify_config_change" => { if payload == "server" && server_mode { tracing::error!( "Server config change detected but server config is obsolete: {}", payload ); } else if worker_mode && payload == format!("worker__{}", *WORKER_GROUP) { tracing::info!("Worker config change detected: {}", payload); reload_worker_config(db, tx.clone(), true).await; } else { tracing::debug!("config changed but did not target this server/worker"); } } "restart_worker_group" => { if worker_mode && payload == *WORKER_GROUP { tracing::info!("Restart requested for worker group '{payload}'"); spawn_graceful_killpill(tx, db, 30, "worker group restart requested", server_mode) .await; } } "notify_webhook_change" => { tracing::info!( "Webhook change detected, invalidating webhook cache: {}", payload ); windmill_api::webhook_util::WEBHOOK_CACHE.remove(payload); } "notify_workspace_envs_change" => { tracing::info!( "Workspace envs change detected, invalidating workspace envs cache: {}", payload ); windmill_common::variables::CUSTOM_ENVS_CACHE.remove(payload); } "notify_asset_producer_change" => { tracing::debug!( "Asset producer change for workspace {}, invalidating producer-writes cache", payload ); windmill_queue::asset_dispatch::ASSET_PRODUCER_WRITES_CACHE.remove(payload); } "notify_macro_registry_change" => { tracing::debug!( "Macro registry change for workspace {}, invalidating macro registry cache", payload ); windmill_common::assets::MACRO_REGISTRY_CACHE.remove(payload); } "notify_workspace_key_change" => { tracing::info!( "Workspace key change detected, invalidating workspace key cache: {}", payload ); windmill_common::variables::WORKSPACE_CRYPT_CACHE.remove(payload); } c if c == windmill_queue::tags::FORK_LINEAGE_CHANGE_CHANNEL => { tracing::info!( "Fork lineage change detected ({payload}), dropping lineage-derived caches" ); windmill_queue::tags::apply_fork_lineage_change(payload); } "notify_workspace_premium_change" => { tracing::info!( "Workspace premium change detected, invalidating workspace premium cache: {}", payload ); windmill_common::workspaces::TEAM_PLAN_CACHE.remove(payload); } "notify_workspace_rate_limit_change" => { tracing::info!( "Workspace rate limit change detected, invalidating rate limit cache: {}", payload ); windmill_common::workspaces::PUBLIC_APP_RATE_LIMIT_CACHE.remove(payload); } "notify_runnable_version_change" => { tracing::info!("Runnable version change detected: {}", payload); match payload.split(':').collect::>().as_slice() { [workspace_id, source_type, path, kind] => { let key = (workspace_id.to_string(), path.to_string()); match *source_type { "script" => { // Evicts DEPLOYED_SCRIPT_HASH_CACHE together with the bundle-cache // key resolution for imported scripts (IMPORTED_SCRIPT_HASH_CACHE), // so key and inlined content flip to the new version in the same // window as the content-side cache below. windmill_common::invalidate_latest_script_hash_caches( workspace_id, path, ); // Evict the relative-import latest-hash cache so a redeployed // imported script flips the content cache to its new version // across all replicas within a poll interval (see #6769). Keyed // by the bare path, matching this event's payload. windmill_api_scripts::scripts::RAW_SCRIPT_LATEST_HASH_CACHE .remove(&format!("{workspace_id}:{path}")); if *kind == "preprocessor" { match sqlx::query_scalar::<_, i64>( "SELECT fv.id FROM flow f INNER JOIN flow_version fv ON fv.id = f.versions[array_upper(f.versions, 1)] WHERE fv.value->'preprocessor_module'->'value'->>'path' = $1 AND f.workspace_id = $2", ) .bind(*path) .bind(*workspace_id) .fetch_all(db).await { Ok(flow_versions) => { tracing::debug!("Workspace preprocessor {} changed, removing runnable format version cache for flow versions {:?}", path, flow_versions); for version in flow_versions { for trigger_kind in TriggerKind::iter() { let key = (windmill_common::triggers::HubOrWorkspaceId::WorkspaceId(workspace_id.to_string()), version, trigger_kind); windmill_common::triggers::RUNNABLE_FORMAT_VERSION_CACHE.remove(&key); } } } Err(e) => { tracing::error!("Error fetching flow paths: {e:#}"); } } } } "flow" => { let dynamic_input_key = windmill_common::jobs::generate_dynamic_input_key( workspace_id, path, ); windmill_common::DYNAMIC_INPUT_CACHE.remove(&dynamic_input_key); windmill_common::FLOW_VERSION_CACHE.remove(&key); } _ => { tracing::warn!("Unknown runnable version change payload: {}", payload); } } } _ => { tracing::warn!("Unknown runnable version change payload: {}", payload); } } } #[cfg(feature = "http_trigger")] "notify_http_trigger_change" => { tracing::info!("HTTP trigger change detected: {}", payload); match windmill_api::triggers::http::refresh_routers(db).await { Ok((true, _)) => { tracing::info!("Refreshed HTTP routers (trigger change)"); } Ok((false, _)) => { tracing::warn!( "Should have refreshed HTTP routers (trigger change) but did not" ); } Err(err) => { tracing::error!("Error refreshing HTTP routers (trigger change): {err:#}"); } }; } "notify_token_invalidation" => { tracing::info!( "Token invalidation detected for prefix: {}...", payload.get(..8).unwrap_or(payload) ); windmill_api::auth::invalidate_token_from_cache(payload); } "notify_app_policy_change" => { // payload is `:`; workspace ids can't contain ':'. if server_mode { if let Some((workspace_id, path)) = payload.split_once(':') { tracing::info!("App policy change detected, invalidating cache: {payload}"); windmill_api::invalidate_app_policy_cache(workspace_id, path); } } } "notify_global_setting_change" => { tracing::info!("Global setting change detected: {}", payload); match payload { BASE_URL_SETTING => { if let Err(e) = reload_base_url_setting(conn).await { tracing::error!(error = %e, "Could not reload base url setting"); } } OAUTH_SETTING => { if let Err(e) = reload_base_url_setting(conn).await { tracing::error!(error = %e, "Could not reload oauth setting"); } } CUSTOM_TAGS_SETTING => { if let Err(e) = reload_custom_tags_setting(db).await { tracing::error!(error = %e, "Could not reload custom tags setting"); } } LICENSE_KEY_SETTING => { if let Err(e) = reload_license_key(&db.into()).await { tracing::error!("Failed to reload license key: {e:#}"); } } DEFAULT_TAGS_PER_WORKSPACE_SETTING => { if let Err(e) = load_tag_per_workspace_enabled(db).await { tracing::error!("Error loading default tag per workspace: {e:#}"); } } DEFAULT_TAGS_WORKSPACES_SETTING => { if let Err(e) = load_tag_per_workspace_workspaces(db).await { tracing::error!( "Error loading default tag per workspace workspaces: {e:#}" ); } } FORK_WORKSPACE_TAG_APPEND_FORK_SUFFIX_SETTING => { if let Err(e) = load_fork_workspace_tag_append_fork_suffix(db).await { tracing::error!( "Error loading fork workspace tag append fork suffix: {e:#}" ); } } PREVIEW_TAGS_OVERRIDE_SETTING => { if let Err(e) = load_preview_tags_override(db).await { tracing::error!("Error loading preview tags override: {e:#}"); } } WORKSPACE_FAIRNESS_ENABLED_SETTING => { if let Err(e) = load_workspace_fairness_enabled(db).await { tracing::error!("Error loading workspace fairness enabled: {e:#}"); } } WORKSPACE_FAIRNESS_MAX_PERCENT_SETTING => { if let Err(e) = load_workspace_fairness_max_percent(db).await { tracing::error!("Error loading workspace fairness max percent: {e:#}"); } } WORKSPACE_FAIRNESS_DURATION_SECS_SETTING => { if let Err(e) = load_workspace_fairness_duration_secs(db).await { tracing::error!("Error loading workspace fairness duration secs: {e:#}"); } } WORKSPACE_FAIRNESS_MIN_TOTAL_SETTING => { if let Err(e) = load_workspace_fairness_min_total(db).await { tracing::error!("Error loading workspace fairness min total: {e:#}"); } } CONCURRENCY_KEY_MAX_QUEUED_SETTING => { if let Err(e) = load_concurrency_key_max_queued(db).await { tracing::error!("Error loading concurrency key max queued: {e:#}"); } } WORKSPACE_MAX_QUEUED_JOBS_SETTING => { if let Err(e) = load_workspace_max_queued_jobs(db).await { tracing::error!("Error loading workspace max queued jobs: {e:#}"); } } SMTP_SETTING => { reload_smtp_config(db).await; } TEAMS_SETTING => { tracing::info!("Teams setting changed."); } INDEXER_SETTING => { reload_indexer_config(db).await; } TIMEOUT_WAIT_RESULT_SETTING => reload_timeout_wait_result_setting(conn).await, RETENTION_PERIOD_SECS_SETTING => reload_retention_period_setting(conn).await, RETENTION_PERIOD_SECS_OVERRIDES_SETTING => { if let Err(e) = load_retention_period_overrides(db).await { tracing::error!("Error loading per-workspace retention overrides: {e:#}"); } } AUDIT_LOG_RETENTION_DAYS_SETTING => { reload_audit_log_retention_days_setting(conn).await } MONITOR_LOGS_ON_OBJECT_STORE_SETTING => { reload_delete_logs_periodically_setting(conn).await } STORE_AUDIT_LOGS_S3_SETTING => reload_store_audit_logs_s3_setting(conn).await, JOB_DEFAULT_TIMEOUT_SECS_SETTING => reload_job_default_timeout_setting(conn).await, JOB_ISOLATION_SETTING => reload_job_isolation_setting(conn).await, NSJAIL_TMPFS_SIZE_MB_SETTING => reload_nsjail_tmpfs_size_setting(conn).await, NSJAIL_TMP_BACKING_SETTING => reload_nsjail_tmp_backing_setting(conn).await, SANDBOX_IMAGE_MAX_SIZE_MB_SETTING => { reload_sandbox_image_max_size_setting(conn).await } SANDBOX_IMAGE_CACHE_MAX_MB_SETTING => { reload_sandbox_image_cache_max_setting(conn).await } SANDBOX_IMAGE_PULL_POLICY_SETTING => { reload_sandbox_image_pull_policy_setting(conn).await } SANDBOX_IMAGE_DEFAULT_REGISTRY_SETTING => { reload_sandbox_image_default_registry_setting(conn).await } SANDBOX_REGISTRY_AUTH_SETTING => reload_sandbox_registry_auth_setting(conn).await, #[cfg(feature = "parquet")] OBJECT_STORE_CONFIG_SETTING => { if !disable_s3_store { reload_object_store_setting(db).await; } } SCIM_TOKEN_SETTING => reload_scim_token_setting(conn).await, EXTRA_PIP_INDEX_URL_SETTING => reload_extra_pip_index_url_setting(conn).await, PIP_INDEX_URL_SETTING => reload_pip_index_url_setting(conn).await, UV_INDEX_STRATEGY_SETTING => reload_uv_index_strategy_setting(conn).await, UV_EXCLUDE_NEWER_SETTING => reload_uv_exclude_newer_setting(conn).await, UV_PYTHON_INSTALL_MIRROR_SETTING => { reload_uv_python_install_mirror_setting(conn).await } BUN_INSTALL_MIN_RELEASE_AGE_SETTING => { reload_bun_install_min_release_age_setting(conn).await } INSTANCE_PYTHON_VERSION_SETTING => { reload_instance_python_version_setting(conn).await } NPM_CONFIG_REGISTRY_SETTING => reload_npm_config_registry_setting(conn).await, BUNFIG_INSTALL_SCOPES_SETTING => reload_bunfig_install_scopes_setting(conn).await, NUGET_CONFIG_SETTING => reload_nuget_config_setting(conn).await, POWERSHELL_REPO_URL_SETTING => reload_powershell_repo_url_setting(conn).await, POWERSHELL_REPO_PAT_SETTING => reload_powershell_repo_pat_setting(conn).await, MAVEN_REPOS_SETTING => reload_maven_repos_setting(conn).await, MAVEN_SETTINGS_XML_SETTING => reload_maven_settings_xml_setting(conn).await, NO_DEFAULT_MAVEN_SETTING => reload_no_default_maven_setting(conn).await, RUBY_REPOS_SETTING => reload_ruby_repos_setting(conn).await, WORKSPACE_REGISTRIES_SETTING => reload_workspace_registries_setting(conn).await, HUB_API_SECRET_SETTING => reload_hub_api_secret_setting(conn).await, KEEP_JOB_DIR_SETTING => { load_keep_job_dir(conn).await; } OTEL_TRACING_PROXY_SETTING => { reload_otel_tracing_proxy_setting(conn).await; if worker_mode { tracing::info!("OTEL tracing proxy setting changed, restarting worker"); spawn_graceful_killpill( tx, db, 30, "OTEL tracing proxy setting change", server_mode, ) .await; } } REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING => { load_require_preexisting_user(db).await; } DISABLE_PASSWORD_LOGIN_SETTING => { load_disable_password_login(db).await; } EXPOSE_METRICS_SETTING => { tracing::info!("Metrics setting changed, restarting"); spawn_graceful_killpill(tx, db, 30, "metrics setting change", server_mode) .await; } EMAIL_DOMAIN_SETTING => { tracing::info!("Email domain setting changed"); if server_mode { spawn_graceful_killpill( tx, db, 30, "email domain setting change", server_mode, ) .await; } } EXPOSE_DEBUG_METRICS_SETTING => { if let Err(e) = load_metrics_debug_enabled(conn).await { tracing::error!(error = %e, "Could not reload debug metrics setting"); } } APP_WORKSPACED_ROUTE_SETTING => { if let Err(e) = reload_app_workspaced_route_setting(db).await { tracing::error!(error = %e, "Could not reload app workspaced route setting"); } } HTTP_ROUTE_WORKSPACED_ROUTE_SETTING => { if let Err(e) = reload_http_route_workspaced_route_setting(db).await { tracing::error!(error = %e, "Could not reload http route workspaced route setting"); } #[cfg(feature = "http_trigger")] match windmill_api::triggers::http::refresh_routers(db).await { Ok((true, _)) => { tracing::info!( "Refreshed HTTP routers (http workspaced route setting change)" ); } Err(err) => { tracing::error!("Error refreshing HTTP routers (http workspaced route setting change): {err:#}"); } _ => {} } } AI_CONFIG_SETTING => { tracing::info!("AI config setting changed, bumping instance AI cache revision"); bump_instance_ai_config_revision(); } OTEL_SETTING => { tracing::info!("OTEL setting changed, restarting"); spawn_graceful_killpill(tx, db, 30, "OTEL setting change", server_mode).await; } REQUEST_SIZE_LIMIT_SETTING => { if server_mode { tracing::info!("Request limit size change detected, killing server expecting to be restarted"); spawn_graceful_killpill( tx, db, 30, "request size limit change", server_mode, ) .await; } } SAML_METADATA_SETTING => { tracing::info!( "SAML metadata change detected, killing server expecting to be restarted" ); spawn_graceful_killpill(tx, db, 30, "SAML metadata change", server_mode).await; } HUB_BASE_URL_SETTING => { if let Err(e) = reload_hub_base_url_setting(conn, server_mode).await { tracing::error!(error = %e, "Could not reload hub base url setting"); } } CRITICAL_ERROR_CHANNELS_SETTING => { if let Err(e) = reload_critical_error_channels_setting(db).await { tracing::error!(error = %e, "Could not reload critical error emails setting"); } } CRITICAL_ALERTS_ON_DB_OVERSIZE_SETTING => { if let Err(e) = reload_critical_alerts_on_db_oversize(db).await { tracing::error!(error = %e, "Could not reload critical alerts on db oversize setting"); } } JWT_SECRET_SETTING => { if let Err(e) = reload_jwt_secret_setting(db).await { tracing::error!(error = %e, "Could not reload jwt secret setting"); } } CRITICAL_ALERT_MUTE_UI_SETTING => { tracing::info!("Critical alert UI setting changed"); if let Err(e) = reload_critical_alert_mute_ui_setting(conn).await { tracing::error!(error = %e, "Could not reload critical alert UI setting"); } } CRITICAL_ALERTS_ON_TOKEN_EXPIRY_SETTING => { if let Err(e) = reload_critical_alerts_on_token_expiry_setting(conn).await { tracing::error!(error = %e, "Could not reload critical alerts on token expiry setting"); } } INSTANCE_EVENTS_WEBHOOK_SETTING => { reload_instance_events_webhook_setting(db).await; } "workspace_telemetry_enabled" => { // Read the new value from the database and log it let enabled = sqlx::query_scalar!( "SELECT value FROM global_settings WHERE name = 'workspace_telemetry_enabled'" ) .fetch_optional(db) .await .ok() .flatten() .and_then(|v| v.as_bool()) .unwrap_or(false); tracing::info!("Workspace telemetry setting changed: enabled={}", enabled); } RESTART_COORDINATION_SETTING => { // Internal coordination key for staggered restarts, no action needed } "plain_emails_telemetry" => { let enabled = sqlx::query_scalar!( "SELECT value FROM global_settings WHERE name = 'plain_emails_telemetry'" ) .fetch_optional(db) .await .ok() .flatten() .and_then(|v| v.as_bool()) .unwrap_or(false); tracing::info!( "Plain emails telemetry setting changed: enabled={}", enabled ); } _ => { tracing::info!("Unrecognized Global Setting Change Payload: {:?}", payload); } } } _ => { tracing::warn!("Unknown notification channel: {}", channel); } } } fn display_config(envs: &[&str]) { tracing::info!( "config: {}", envs.iter() .filter(|env| std::env::var(env).is_ok()) .map(|env| { format!( "{}: {}", env, std::env::var(env).unwrap_or_else(|_| "not set".to_string()) ) }) .collect::>() .join(", ") ) } pub struct WorkerConn { conn: Connection, worker_name: String, } pub async fn run_workers( mut rx: tokio::sync::broadcast::Receiver<()>, tx: KillpillSender, base_internal_url: String, hostname: String, workers: &[WorkerConn], ) -> anyhow::Result<()> { let mut killpill_rxs = vec![]; let num_workers = workers.len(); for _ in 0..num_workers { killpill_rxs.push(rx.resubscribe()); } if rx.try_recv().is_ok() { tracing::info!("Received killpill, exiting"); return Ok(()); } // #[cfg(tokio_unstable)] // let monitor = tokio_metrics::TaskMonitor::new(); windmill_common::external_ip::resolve_ip_in_background(); let mut handles = Vec::with_capacity(num_workers as usize); for x in [ &*TMP_LOGS_DIR, &*UV_CACHE_DIR, &*DENO_CACHE_DIR, &*DENO_CACHE_DIR_DEPS, &*DENO_CACHE_DIR_NPM, &*BUN_CACHE_DIR, &*PY310_CACHE_DIR, &*PY311_CACHE_DIR, &*PY312_CACHE_DIR, &*PY313_CACHE_DIR, &*BUN_BUNDLE_CACHE_DIR, &*GO_CACHE_DIR, &*GO_BIN_CACHE_DIR, &*RUST_CACHE_DIR, &*CSHARP_CACHE_DIR, &*NU_CACHE_DIR, &*HUB_CACHE_DIR, &*POWERSHELL_CACHE_DIR, &*JAVA_CACHE_DIR, &*RUBY_CACHE_DIR, &*R_CACHE_DIR, &*TAR_JAVA_CACHE_DIR, // for related places search: ADD_NEW_LANG ] { DirBuilder::new() .recursive(true) .create(x) .expect("could not create initial worker dir"); } tracing::info!( "Starting {num_workers} workers and SLEEP_QUEUE={}ms", windmill_worker::sleep_queue() ); for i in 1..(num_workers + 1) { let wk_conf = &workers[i as usize - 1]; let conn1 = wk_conf.conn.clone(); let worker_name = wk_conf.worker_name.clone(); WORKERS_NAMES.write().await.push(worker_name.clone()); let rx = killpill_rxs.pop().unwrap(); let tx = tx.clone(); let base_internal_url = base_internal_url.clone(); let hostname = hostname.clone(); handles.push(tokio::spawn(async move { if num_workers > 1 { tracing::info!(worker = %worker_name, "starting worker {i}"); } let f = windmill_worker::run_worker( &conn1, &hostname, worker_name, i as u64, num_workers as u32, rx, tx, &base_internal_url, ); // #[cfg(tokio_unstable)] // { // monitor.monitor(f, "worker").await // } // #[cfg(not(tokio_unstable))] // { f.await // } })); } futures::future::try_join_all(handles).await?; Ok(()) } /// Schedule a graceful restart with DB-coordinated staggering. /// /// Uses a PostgreSQL advisory lock to serialize restart scheduling across server instances. /// Each instance records its planned restart time in the `_restart_coordination` global setting; /// subsequent instances read existing schedules and shift their restart to maintain at least /// `safety_margin_secs` between consecutive restarts (must exceed the server startup time). /// /// Every server waits at least `DRAIN_DELAY_SECS` to let in-flight requests complete. /// Each subsequent server waits an additional `safety_margin_secs` after the previous one, /// guaranteeing zero downtime overlap. /// /// The DB coordination is done synchronously (fast, ~ms) to reserve our restart slot, /// then the sleep+kill is spawned in the background so the notification handler is not blocked. /// /// Falls back to drain-only delay if DB coordination fails. /// /// Only `server_mode` processes coordinate, on the strength of the worker case: a worker /// group restarting costs queue latency rather than lost work, `v2_job_queue` being durable. /// Were workers to take part, one could claim the `is_first` slot and leave every server /// holding its shutdown open for a peer that serves no API traffic. async fn spawn_graceful_killpill( tx: &KillpillSender, db: &Pool, safety_margin_secs: u64, context: &str, server_mode: bool, ) { // Minimum delay before any restart to let in-flight requests drain const DRAIN_DELAY_SECS: u64 = 3; let (delay, is_first) = if !server_mode { (DRAIN_DELAY_SECS, true) } else { match coordinate_restart_delay(db, safety_margin_secs, DRAIN_DELAY_SECS).await { Ok(r) => r, Err(e) => { tracing::warn!( "Failed to coordinate restart for {context}: {e:#}, \ falling back to drain delay of {DRAIN_DELAY_SECS}s" ); (DRAIN_DELAY_SECS, true) } } }; tracing::info!( "Scheduling {context} graceful shutdown in {delay}s (first_to_restart={is_first})" ); let tx = tx.clone(); let db = db.clone(); let context = context.to_string(); // Capture the current time so the health check only considers heartbeats // written *after* this restart was initiated (not stale pre-restart ones). let not_before = chrono::Utc::now(); tokio::spawn(async move { if !is_first { // Non-first instance: wait for another server to announce itself as // started (via the coordination record) before shutting down, ensuring // zero-downtime. Falls back to the computed delay as a maximum timeout. let started_waiting = tokio::time::Instant::now(); let max_wait = Duration::from_secs(delay.max(safety_margin_secs).max(120)); loop { if windmill_api::check_any_server_started(&db, not_before).await { tracing::info!( "{context}: healthy peer detected after {:.1}s, proceeding with shutdown", started_waiting.elapsed().as_secs_f64() ); break; } if started_waiting.elapsed() >= max_wait { tracing::warn!( "{context}: no healthy peer detected after {:.1}s, proceeding with shutdown anyway", started_waiting.elapsed().as_secs_f64() ); break; } tokio::time::sleep(Duration::from_secs(5)).await; } } else { tokio::time::sleep(Duration::from_secs(delay)).await; } tx.send(); }); } /// Coordinate a restart delay with other instances via the DB. /// /// Returns `(delay_secs, is_first)`: /// - `delay_secs`: how long to wait before restarting /// - `is_first`: whether this instance is the first to schedule a restart /// /// The first server restarts immediately (0 delay) to maximize its head-start. /// Subsequent servers use a safety margin after the latest scheduled restart as a /// fallback timeout, but primarily wait for a health-check (see `spawn_graceful_killpill`). async fn coordinate_restart_delay( db: &Pool, safety_margin_secs: u64, drain_delay_secs: u64, ) -> anyhow::Result<(u64, bool)> { const RESTART_LOCK_ID: i64 = 737_483_920; // Stale threshold: ignore coordination entries older than this const STALE_THRESHOLD_SECS: i64 = 120; let now = chrono::Utc::now(); let mut tx = db.begin().await.context("begin restart coordination tx")?; // Serialize access across all instances sqlx::query("SELECT pg_advisory_xact_lock($1)") .bind(RESTART_LOCK_ID) .execute(&mut *tx) .await .context("acquire restart coordination lock")?; // Read existing coordination record let existing: Option = sqlx::query_scalar("SELECT value FROM global_settings WHERE name = $1") .bind(RESTART_COORDINATION_SETTING) .fetch_optional(&mut *tx) .await .context("read restart coordination")?; // Parse existing scheduled restarts, filtering out stale entries // Each entry is (instance_name, restart_at) let mut scheduled: Vec<(String, chrono::DateTime)> = Vec::new(); if let Some(val) = &existing { if let Some(arr) = val.get("restarts").and_then(|v| v.as_array()) { for entry in arr { let instance = entry .get("instance") .and_then(|v| v.as_str()) .unwrap_or("unknown") .to_string(); if let Some(ts_str) = entry.get("restart_at").and_then(|v| v.as_str()) { if let Ok(dt) = chrono::DateTime::parse_from_rfc3339(ts_str) { let dt = dt.with_timezone(&chrono::Utc); let stale_cutoff = now - chrono::Duration::seconds(STALE_THRESHOLD_SECS); if dt > stale_cutoff { scheduled.push((instance, dt)); } } } } } } // Find the latest scheduled restart let latest = scheduled.iter().map(|(_, dt)| *dt).max(); let is_first = latest.is_none(); let earliest_allowed = now + chrono::Duration::seconds(drain_delay_secs as i64); // Our restart time: // - First instance (no prior scheduled restarts): brief 1s delay to let in-flight // requests (like the settings save) complete, then restart quickly to maximize // head-start before subsequent instances shut down. // - Subsequent instances: safety_margin after the latest scheduled restart (used as // a fallback timeout; the actual shutdown is gated by a health-check in the caller). let our_restart = match latest { Some(last) => { let after_last = last + chrono::Duration::seconds(safety_margin_secs as i64); // Use whichever is later: drain delay or staggered position earliest_allowed.max(after_last) } None => now + chrono::Duration::seconds(1), }; // Record our restart time (deduplicate: remove any prior entry for this instance) scheduled.retain(|(inst, _)| inst != &*INSTANCE_NAME); scheduled.push((INSTANCE_NAME.clone(), our_restart)); let new_value = serde_json::json!({ "restarts": scheduled.iter().map(|(inst, dt)| { serde_json::json!({ "instance": inst, "restart_at": dt.to_rfc3339() }) }).collect::>() }); sqlx::query( "INSERT INTO global_settings (name, value, updated_at) \ VALUES ($1, $2, now()) \ ON CONFLICT (name) DO UPDATE SET value = $2, updated_at = now()", ) .bind(RESTART_COORDINATION_SETTING) .bind(&new_value) .execute(&mut *tx) .await .context("write restart coordination")?; tx.commit().await.context("commit restart coordination")?; let delay = (our_restart - now).num_seconds().max(0) as u64; Ok((delay, is_first)) }