diff --git a/.github/workflows/docker-image.yml b/.github/workflows/docker-image.yml index 42cc21264b..f1446210cc 100644 --- a/.github/workflows/docker-image.yml +++ b/.github/workflows/docker-image.yml @@ -138,7 +138,7 @@ jobs: platforms: linux/amd64,linux/arm64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka,otel tags: | ${{ env.REGISTRY }}/${{ env.IMAGE_NAME }}-ee:${{ env.DEV_SHA }} ${{ steps.meta-ee-public.outputs.tags }} @@ -200,7 +200,7 @@ jobs: platforms: linux/amd64 push: true build-args: | - features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka + features=enterprise,enterprise_saml,stripe,embedding,parquet,prometheus,openidconnect,cloud,jemalloc,tantivy,deno_core,kafka,otel PYTHON_IMAGE=python:3.12.2-slim-bookworm tags: | ${{ steps.meta-ee-public-py312.outputs.tags }} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index dac85d4f09..e003c2300a 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -4173,6 +4173,19 @@ dependencies = [ "tower-service", ] +[[package]] +name = "hyper-timeout" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3203a961e5c83b6f5498933e78b6b263e208c197b63e9c6c53cc82ffd3f63793" +dependencies = [ + "hyper 1.5.1", + "hyper-util", + "pin-project-lite", + "tokio", + "tower-service", +] + [[package]] name = "hyper-tls" version = "0.5.0" @@ -4809,7 +4822,7 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "674883a98273598ac3aad4301724c56734bea90574c5033af067e8f9fb5eb399" dependencies = [ - "prost", + "prost 0.12.6", "prost-types", ] @@ -5637,6 +5650,91 @@ dependencies = [ "vcpkg", ] +[[package]] +name = "opentelemetry" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0f3cebff57f7dbd1255b44d8bddc2cebeb0ea677dbaa2e25a3070a91b318f660" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "once_cell", + "pin-project-lite", + "thiserror 1.0.69", +] + +[[package]] +name = "opentelemetry-appender-tracing" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ab5feffc321035ad94088a7e5333abb4d84a8726e54a802e736ce9dd7237e85b" +dependencies = [ + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", +] + +[[package]] +name = "opentelemetry-otlp" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91cf61a1868dacc576bf2b2a1c3e9ab150af7272909e80085c3173384fe11f76" +dependencies = [ + "async-trait", + "futures-core", + "http 1.2.0", + "opentelemetry", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost 0.13.3", + "thiserror 1.0.69", + "tokio", + "tonic", + "tracing", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a6e05acbfada5ec79023c85368af14abd0b307c015e9064d249b2a950ef459a6" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost 0.13.3", + "tonic", +] + +[[package]] +name = "opentelemetry-semantic-conventions" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bc1b6902ff63b32ef6c489e8048c5e253e2e4a803ea3ea7e783914536eb15c52" + +[[package]] +name = "opentelemetry_sdk" +version = "0.27.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "27b742c1cae4693792cc564e58d75a2a0ba29421a34a85b50da92efa89ecb2bc" +dependencies = [ + "async-trait", + "futures-channel", + "futures-executor", + "futures-util", + "glob", + "once_cell", + "opentelemetry", + "percent-encoding", + "rand 0.8.5", + "serde_json", + "thiserror 1.0.69", + "tokio", + "tokio-stream", + "tracing", +] + [[package]] name = "option-ext" version = "0.2.0" @@ -6264,7 +6362,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "deb1435c188b76130da55f17a466d252ff7b1418b2ad3e037d127b94e3411f29" dependencies = [ "bytes", - "prost-derive", + "prost-derive 0.12.6", +] + +[[package]] +name = "prost" +version = "0.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b0487d90e047de87f984913713b85c601c05609aad5b0df4b4573fbf69aa13f" +dependencies = [ + "bytes", + "prost-derive 0.13.3", ] [[package]] @@ -6280,13 +6388,26 @@ dependencies = [ "syn 2.0.90", ] +[[package]] +name = "prost-derive" +version = "0.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e9552f850d5f0964a4e4d0bf306459ac29323ddfbae05e35a7c0d35cb0803cc5" +dependencies = [ + "anyhow", + "itertools 0.13.0", + "proc-macro2", + "quote", + "syn 2.0.90", +] + [[package]] name = "prost-types" version = "0.12.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9091c90b0a32608e984ff2fa4091273cbdd755d54935c51d520887f4a1dbd5b0" dependencies = [ - "prost", + "prost 0.12.6", ] [[package]] @@ -9448,6 +9569,36 @@ dependencies = [ "winnow 0.6.20", ] +[[package]] +name = "tonic" +version = "0.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "877c5b330756d856ffcc4553ab34a5684481ade925ecc54bcd1bf02b1d0d4d52" +dependencies = [ + "async-stream", + "async-trait", + "axum", + "base64 0.22.1", + "bytes", + "h2 0.4.7", + "http 1.2.0", + "http-body 1.0.1", + "http-body-util", + "hyper 1.5.1", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "prost 0.13.3", + "socket2", + "tokio", + "tokio-stream", + "tower 0.4.13", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "toolchain_find" version = "0.4.0" @@ -9469,11 +9620,16 @@ checksum = "b8fa9be0de6cf49e536ce1851f987bd21a43b771b09473c3549a6c853db37c1c" dependencies = [ "futures-core", "futures-util", + "indexmap 1.9.3", "pin-project", "pin-project-lite", + "rand 0.8.5", + "slab", "tokio", + "tokio-util", "tower-layer", "tower-service", + "tracing", ] [[package]] @@ -9651,6 +9807,24 @@ dependencies = [ "url", ] +[[package]] +name = "tracing-opentelemetry" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97a971f6058498b5c0f1affa23e7ea202057a7301dbff68e968b2d578bcbd053" +dependencies = [ + "js-sys", + "once_cell", + "opentelemetry", + "opentelemetry_sdk", + "smallvec", + "tracing", + "tracing-core", + "tracing-log 0.2.0", + "tracing-subscriber", + "web-time", +] + [[package]] name = "tracing-serde" version = "0.1.3" @@ -10644,6 +10818,11 @@ dependencies = [ "magic-crypt", "mail-send", "object_store", + "opentelemetry", + "opentelemetry-appender-tracing", + "opentelemetry-otlp", + "opentelemetry-semantic-conventions", + "opentelemetry_sdk", "pin-project-lite", "prometheus", "quick_cache", @@ -10662,6 +10841,7 @@ dependencies = [ "tracing-appender", "tracing-flame", "tracing-loki", + "tracing-opentelemetry", "tracing-subscriber", "uuid 1.11.0", "windmill-macros", @@ -10897,6 +11077,7 @@ dependencies = [ "hmac", "itertools 0.13.0", "lazy_static", + "opentelemetry", "prometheus", "regex", "reqwest 0.12.9", @@ -10961,6 +11142,7 @@ dependencies = [ "object_store", "once_cell", "openidconnect", + "opentelemetry", "pem 3.0.4", "postgres-native-tls", "prometheus", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 3f591dfd84..4646af9f5e 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -63,6 +63,7 @@ tantivy = ["dep:windmill-indexer", "windmill-api/tantivy"] sqlx = ["windmill-worker/sqlx"] deno_core = ["windmill-worker/deno_core", "dep:deno_core"] kafka = ["windmill-api/kafka"] +otel = ["windmill-common/otel", "windmill-worker/otel"] [dependencies] anyhow.workspace = true @@ -95,7 +96,6 @@ deno_core = { workspace = true, optional = true } object_store = { workspace = true, optional = true } quote.workspace = true - [target.'cfg(not(target_env = "msvc"))'.dependencies] tikv-jemallocator = { optional = true, workspace = true } tikv-jemalloc-sys = { optional = true, workspace = true } @@ -279,6 +279,13 @@ tar = "^0" http = "^1" async-stream = "^0" +opentelemetry = "0.27.0" +tracing-opentelemetry = "0.28.0" +opentelemetry_sdk = { version = "*", features = ["rt-tokio"] } +opentelemetry-otlp = "0.27.0" +opentelemetry-appender-tracing = "0.27.0" +opentelemetry-semantic-conventions = "*" + tikv-jemallocator = { version = "0.5" } tikv-jemalloc-sys = { version = "^0.5" } tikv-jemalloc-ctl = { version = "^0.5" } diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 6bb5c4686f..b5181c39f2 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -5066da602260334767186e69ae5b6821feca0c71 \ No newline at end of file +ad89b5a1566159490eceb59c3e5c9416d9aa855c \ No newline at end of file diff --git a/backend/src/main.rs b/backend/src/main.rs index d403ad212a..80edff1ac9 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -8,7 +8,7 @@ use anyhow::Context; use monitor::{ - reload_delete_logs_periodically_setting, reload_indexer_config, + load_otel, reload_delete_logs_periodically_setting, reload_indexer_config, reload_timeout_wait_result_setting, send_current_log_file_to_object_store, send_logs_to_object_store, }; @@ -37,7 +37,7 @@ use windmill_common::{ EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING, HUB_BASE_URL_SETTING, INDEXER_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING, JWT_SECRET_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NPM_CONFIG_REGISTRY_SETTING, - OAUTH_SETTING, PIP_INDEX_URL_SETTING, REQUEST_SIZE_LIMIT_SETTING, + OAUTH_SETTING, OTEL_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, SMTP_SETTING, TIMEOUT_WAIT_RESULT_SETTING, }, @@ -221,11 +221,66 @@ async fn windmill_main() -> anyhow::Result<()> { let hostname = hostname(); - #[cfg(not(feature = "flamegraph"))] - let _guard = windmill_common::tracing_init::initialize_tracing(&hostname); + let mut enable_standalone_indexer: bool = false; + + let mode = std::env::var("MODE") + .map(|x| x.to_lowercase()) + .map(|x| { + if &x == "server" { + println!("Binary is in 'server' mode"); + Mode::Server + } else if &x == "worker" { + tracing::info!("Binary is in 'worker' mode"); + #[cfg(windows)] + { + println!("It is highly recommended to use the agent mode instead on windows (MODE=agent) and to pass a BASE_INTERNAL_URL"); + } + Mode::Worker + } else if &x == "agent" { + println!("Binary is in 'agent' mode"); + if std::env::var("BASE_INTERNAL_URL").is_err() { + panic!("BASE_INTERNAL_URL is required in agent mode") + } + if std::env::var("JOB_TOKEN").is_err() { + println!("JOB_TOKEN is not passed, hence workers will still need to create permissions for each job and the DATABASE_URL needs to be of a role that can INSERT into the job_perms table") + } + + #[cfg(not(feature = "enterprise"))] + { + panic!("Agent mode is only available in the EE, ignoring..."); + } + #[cfg(feature = "enterprise")] + Mode::Agent + } else if &x == "indexer" { + tracing::info!("Binary is in 'indexer' mode"); + #[cfg(not(feature = "tantivy"))] + { + eprintln!("Cannot start the indexer because tantivy is not included in this binary/image. Make sure you are using the EE image if you want to access the full text search features."); + panic!("Indexer mode requires compiling with the tantivy feature flag."); + } + #[cfg(feature = "tantivy")] + Mode::Indexer + } else if &x == "standalone+search"{ + enable_standalone_indexer = true; + println!("Binary is in 'standalone' mode with search enabled"); + Mode::Standalone + } + else { + if &x != "standalone" { + eprintln!("mode not recognized, defaulting to standalone: {x}"); + } else { + println!("Binary is in 'standalone' mode"); + } + Mode::Standalone + } + }) + .unwrap_or_else(|_| { + tracing::info!("Mode not specified, defaulting to standalone"); + Mode::Standalone + }); #[cfg(all(not(target_env = "msvc"), feature = "jemalloc"))] - tracing::info!("jemalloc enabled"); + println!("jemalloc enabled"); #[cfg(feature = "flamegraph")] let _guard = windmill_common::tracing_init::setup_flamegraph(); @@ -236,13 +291,13 @@ async fn windmill_main() -> anyhow::Result<()> { "cache" => { #[cfg(feature = "embedding")] { - tracing::info!("Caching embedding model..."); + println!("Caching embedding model..."); windmill_api::embeddings::ModelInstance::load_model_files().await?; - tracing::info!("Cached embedding model"); + println!("Cached embedding model"); } #[cfg(not(feature = "embedding"))] { - tracing::warn!("Embeddings are not enabled, ignoring..."); + println!("Embeddings are not enabled, ignoring..."); } cache_hub_scripts(std::env::args().nth(2)).await?; @@ -256,64 +311,6 @@ async fn windmill_main() -> anyhow::Result<()> { _ => {} } - let mut enable_standalone_indexer: bool = false; - - let mode = std::env::var("MODE") - .map(|x| x.to_lowercase()) - .map(|x| { - if &x == "server" { - tracing::info!("Binary is in 'server' mode"); - Mode::Server - } else if &x == "worker" { - tracing::info!("Binary is in 'worker' mode"); - #[cfg(windows)] - { - tracing::warn!("It is highly recommended to use the agent mode instead on windows (MODE=agent) and to pass a BASE_INTERNAL_URL"); - } - Mode::Worker - } else if &x == "agent" { - tracing::info!("Binary is in 'agent' mode"); - if std::env::var("BASE_INTERNAL_URL").is_err() { - panic!("BASE_INTERNAL_URL is required in agent mode") - } - if std::env::var("JOB_TOKEN").is_err() { - tracing::warn!("JOB_TOKEN is not passed, hence workers will still need to create permissions for each job and the DATABASE_URL needs to be of a role that can INSERT into the job_perms table") - } - - #[cfg(not(feature = "enterprise"))] - { - panic!("Agent mode is only available in the EE, ignoring..."); - } - #[cfg(feature = "enterprise")] - Mode::Agent - } else if &x == "indexer" { - tracing::info!("Binary is in 'indexer' mode"); - #[cfg(not(feature = "tantivy"))] - { - tracing::error!("Cannot start the indexer because tantivy is not included in this binary/image. Make sure you are using the EE image if you want to access the full text search features."); - panic!("Indexer mode requires compiling with the tantivy feature flag."); - } - #[cfg(feature = "tantivy")] - Mode::Indexer - } else if &x == "standalone+search"{ - enable_standalone_indexer = true; - tracing::info!("Binary is in 'standalone' mode with search enabled"); - Mode::Standalone - } - else { - if &x != "standalone" { - tracing::error!("mode not recognized, defaulting to standalone: {x}"); - } else { - tracing::info!("Binary is in 'standalone' mode"); - } - Mode::Standalone - } - }) - .unwrap_or_else(|_| { - tracing::info!("Mode not specified, defaulting to standalone"); - Mode::Standalone - }); - #[allow(unused_mut)] let mut num_workers = if mode == Mode::Server || mode == Mode::Indexer { 0 @@ -325,7 +322,7 @@ async fn windmill_main() -> anyhow::Result<()> { }; if num_workers > 1 { - tracing::warn!( + println!( "We STRONGLY recommend using at most 1 worker per container, use at your own risks" ); } @@ -347,10 +344,15 @@ async fn windmill_main() -> anyhow::Result<()> { IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)) }; - tracing::info!("Connecting to database..."); + println!("Connecting to database..."); let db = windmill_common::connect_db(server_mode, indexer_mode).await?; + + load_otel(&db).await; tracing::info!("Database connected"); + #[cfg(not(feature = "flamegraph"))] + let _guard = windmill_common::tracing_init::initialize_tracing(&hostname, &mode); + let num_version = sqlx::query_scalar!("SELECT version()").fetch_one(&db).await; tracing::info!( @@ -776,6 +778,15 @@ Windmill Community Edition {GIT_VERSION} tracing::error!(error = %e, "Could not reload debug metrics setting"); } }, + OTEL_SETTING => { + tracing::info!("OTEL setting changed, restarting"); + // we wait a bit randomly to avoid having all servers and workers shutdown at same time + let rd_delay = rand::thread_rng().gen_range(0..4); + tokio::time::sleep(Duration::from_secs(rd_delay)).await; + if let Err(e) = tx.send(()) { + tracing::error!(error = %e, "Could not send killpill"); + } + }, REQUEST_SIZE_LIMIT_SETTING => { if server_mode { tracing::info!("Request limit size change detected, killing server expecting to be restarted"); diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index df3b695168..262927cce3 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -12,7 +12,7 @@ use std::{ use chrono::{NaiveDateTime, Utc}; use futures::{stream::FuturesUnordered, StreamExt}; -use serde::de::DeserializeOwned; +use serde::{de::DeserializeOwned, Deserializer}; use sqlx::{Pool, Postgres}; use tokio::{ join, @@ -40,7 +40,7 @@ use windmill_common::{ EXTRA_PIP_INDEX_URL_SETTING, HUB_BASE_URL_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING, JWT_SECRET_SETTING, KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NPM_CONFIG_REGISTRY_SETTING, OAUTH_SETTING, - PIP_INDEX_URL_SETTING, REQUEST_SIZE_LIMIT_SETTING, + OTEL_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, TIMEOUT_WAIT_RESULT_SETTING, }, @@ -58,7 +58,8 @@ use windmill_common::{ }, BASE_URL, CRITICAL_ALERT_MUTE_UI_ENABLED, CRITICAL_ERROR_CHANNELS, DB, DEFAULT_HUB_BASE_URL, HUB_BASE_URL, JOB_RETENTION_SECS, METRICS_DEBUG_ENABLED, METRICS_ENABLED, - MONITOR_LOGS_ON_OBJECT_STORE, SERVICE_LOG_RETENTION_SECS, + MONITOR_LOGS_ON_OBJECT_STORE, OTEL_LOGS_ENABLED, OTEL_METRICS_ENABLED, OTEL_TRACING_ENABLED, + SERVICE_LOG_RETENTION_SECS, }; use windmill_queue::cancel_job; use windmill_worker::{ @@ -199,6 +200,65 @@ pub async fn load_metrics_enabled(db: &DB) -> error::Result<()> { Ok(()) } +fn empty_string_as_none<'de, D>(deserializer: D) -> Result, D::Error> +where + D: Deserializer<'de>, +{ + let option = as serde::Deserialize>::deserialize(deserializer)?; + Ok(option.filter(|s| !s.is_empty())) +} + +#[derive(serde::Deserialize)] +struct OtelSetting { + metrics_enabled: Option, + logs_enabled: Option, + tracing_enabled: Option, + #[serde(default, deserialize_with = "empty_string_as_none")] + otel_exporter_otlp_endpoint: Option, + #[serde(default, deserialize_with = "empty_string_as_none")] + otel_exporter_otlp_headers: Option, + #[serde(default, deserialize_with = "empty_string_as_none")] + otel_exporter_otlp_protocol: Option, + #[serde(default, deserialize_with = "empty_string_as_none")] + otel_exporter_otlp_compression: Option, +} + +pub async fn load_otel(db: &DB) { + let otel = load_value_from_global_settings(db, OTEL_SETTING).await; + if let Ok(v) = otel { + if let Some(v) = v { + let deser = serde_json::from_value::(v); + if let Ok(o) = deser { + let metrics_enabled = o.metrics_enabled.unwrap_or(false); + let logs_enabled = o.logs_enabled.unwrap_or(false); + let tracing_enabled = o.tracing_enabled.unwrap_or(false); + + OTEL_METRICS_ENABLED.store(metrics_enabled, Ordering::Relaxed); + OTEL_LOGS_ENABLED.store(logs_enabled, Ordering::Relaxed); + OTEL_TRACING_ENABLED.store(tracing_enabled, Ordering::Relaxed); + if let Some(endpoint) = o.otel_exporter_otlp_endpoint.as_ref() { + std::env::set_var("OTEL_EXPORTER_OTLP_ENDPOINT", endpoint); + } + if let Some(headers) = o.otel_exporter_otlp_headers.as_ref() { + std::env::set_var("OTEL_EXPORTER_OTLP_HEADERS", headers); + } + if let Some(protocol) = o.otel_exporter_otlp_protocol { + std::env::set_var("OTEL_EXPORTER_OTLP_PROTOCOL", protocol); + } + if let Some(compression) = o.otel_exporter_otlp_compression { + std::env::set_var("OTEL_EXPORTER_OTLP_COMPRESSION", compression); + } + tracing::info!("OTEL settings loaded: tracing ({tracing_enabled}), logs ({logs_enabled}), metrics ({metrics_enabled}), endpoint ({:?}), headers defined: ({})", + o.otel_exporter_otlp_endpoint, o.otel_exporter_otlp_headers.is_some()); + } else { + tracing::error!("Error deserializing otel settings"); + } + } + } else { + tracing::error!("Error loading otel settings: {}", otel.unwrap_err()); + } +} + pub async fn load_tag_per_workspace_enabled(db: &DB) -> error::Result<()> { let metrics_enabled = load_value_from_global_settings(db, DEFAULT_TAGS_PER_WORKSPACE_SETTING).await; diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index a229450651..2aebf6f5e8 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -77,7 +77,10 @@ async fn initialize_tracing() { static ONCE: Once = Once::new(); ONCE.call_once(|| { - let _ = windmill_common::tracing_init::initialize_tracing("test"); + let _ = windmill_common::tracing_init::initialize_tracing( + "test", + &windmill_common::utils::Mode::Standalone, + ); }); } diff --git a/backend/windmill-api/src/tracing_init.rs b/backend/windmill-api/src/tracing_init.rs index 60dab73810..c12cbe3192 100644 --- a/backend/windmill-api/src/tracing_init.rs +++ b/backend/windmill-api/src/tracing_init.rs @@ -29,11 +29,19 @@ impl OnResponse for MyOnResponse { _span: &tracing::Span, ) { if *LOG_REQUESTS { - tracing::info!( - latency = latency.as_millis(), - status = response.status().as_u16(), - "response" - ) + if response.status().is_success() { + tracing::info!( + latency = latency.as_millis(), + status = response.status().as_u16(), + "response" + ) + } else { + tracing::error!( + latency = latency.as_millis(), + status = response.status().as_u16(), + "response" + ) + } } } } diff --git a/backend/windmill-common/Cargo.toml b/backend/windmill-common/Cargo.toml index b56b9a39cf..5b9a74c361 100644 --- a/backend/windmill-common/Cargo.toml +++ b/backend/windmill-common/Cargo.toml @@ -13,6 +13,7 @@ flamegraph = ["dep:tracing-flame"] loki = ["dep:tracing-loki"] benchmark = [] parquet = ["dep:object_store", "dep:aws-config", "dep:aws-sdk-sts"] +otel = ["dep:opentelemetry-semantic-conventions", "dep:opentelemetry-otlp", "dep:opentelemetry_sdk", "dep:opentelemetry", "dep:tracing-opentelemetry", "dep:opentelemetry-appender-tracing"] [lib] name = "windmill_common" @@ -64,5 +65,12 @@ croner = "2.0.6" quick_cache.workspace = true pin-project-lite.workspace = true +opentelemetry-semantic-conventions = { workspace = true, optional = true } +opentelemetry-otlp = { workspace = true, optional = true } +opentelemetry_sdk = { workspace = true, optional = true } +opentelemetry = { workspace = true, optional = true } +tracing-opentelemetry = { workspace = true, optional = true } +opentelemetry-appender-tracing = { workspace = true, optional = true } + [target.'cfg(not(target_env = "msvc"))'.dependencies] tikv-jemalloc-ctl = { optional = true, workspace = true } diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index 7b932d8f50..9dbfd03b0e 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -35,8 +35,9 @@ pub const CRITICAL_ALERT_MUTE_UI_SETTING: &str = "critical_alert_mute_ui"; pub const DEV_INSTANCE_SETTING: &str = "dev_instance"; pub const JWT_SECRET_SETTING: &str = "jwt_secret"; pub const EMAIL_DOMAIN_SETTING: &str = "email_domain"; +pub const OTEL_SETTING: &str = "otel"; -pub const ENV_SETTINGS: [&str; 51] = [ +pub const ENV_SETTINGS: [&str; 54] = [ "DISABLE_NSJAIL", "MODE", "NUM_WORKERS", @@ -88,4 +89,7 @@ pub const ENV_SETTINGS: [&str; 51] = [ "WORKER_GROUP", "SAML_METADATA", "INSTANCE_IS_DEV", + "OTEL_METRICS", + "OTEL_TRACING", + "OTEL_LOGS", ]; diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index ad47b198af..4140f6e58b 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -36,20 +36,20 @@ pub mod job_s3_helpers_ee; pub mod jobs; pub mod more_serde; pub mod oauth2; +pub mod otel_ee; pub mod queue; pub mod s3_helpers; pub mod schedule; pub mod scripts; pub mod server; pub mod stats_ee; +pub mod tracing_init; pub mod users; pub mod utils; pub mod variables; pub mod worker; pub mod workspaces; -pub mod tracing_init; - pub const DEFAULT_MAX_CONNECTIONS_SERVER: u32 = 50; pub const DEFAULT_MAX_CONNECTIONS_WORKER: u32 = 5; pub const DEFAULT_MAX_CONNECTIONS_INDEXER: u32 = 5; @@ -87,6 +87,12 @@ lazy_static::lazy_static! { .unwrap_or_else(|| SocketAddr::from(([0, 0, 0, 0], *METRICS_PORT))); pub static ref METRICS_ENABLED: AtomicBool = AtomicBool::new(std::env::var("METRICS_PORT").is_ok() || std::env::var("METRICS_ADDR").is_ok()); + + pub static ref OTEL_METRICS_ENABLED: AtomicBool = AtomicBool::new(std::env::var("OTEL_METRICS").is_ok()); + pub static ref OTEL_TRACING_ENABLED: AtomicBool = AtomicBool::new(std::env::var("OTEL_TRACING").is_ok()); + pub static ref OTEL_LOGS_ENABLED: AtomicBool = AtomicBool::new(std::env::var("OTEL_LOGS").is_ok()); + + pub static ref METRICS_DEBUG_ENABLED: AtomicBool = AtomicBool::new(false); pub static ref CRITICAL_ALERT_MUTE_UI_ENABLED: AtomicBool = AtomicBool::new(false); diff --git a/backend/windmill-common/src/otel_ee.rs b/backend/windmill-common/src/otel_ee.rs new file mode 100644 index 0000000000..f310885c0b --- /dev/null +++ b/backend/windmill-common/src/otel_ee.rs @@ -0,0 +1,54 @@ +/* + * 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 crate::{jobs::QueuedJob, utils::Mode}; +use uuid::Uuid; + +pub fn set_span_parent(_span: &tracing::Span, _rj: &Uuid) {} + +#[cfg(not(all(feature = "otel", feature = "enterprise")))] +pub(crate) type OtelProvider = Option<()>; + +#[cfg(all(feature = "otel", feature = "enterprise"))] +pub(crate) type OtelProvider = Option; + +#[cfg(not(feature = "otel"))] +pub fn otel_ctx() -> () {} + +#[cfg(feature = "otel")] +#[inline(always)] +pub fn otel_ctx() -> opentelemetry::Context { + opentelemetry::Context::current() +} + +#[cfg(not(feature = "otel"))] +impl FutureExt for T {} + +#[cfg(not(feature = "otel"))] +pub trait FutureExt: Sized { + fn with_context(self, _otel_cx: ()) -> Self { + self + } +} + +use tracing_subscriber::EnvFilter; + +pub(crate) fn init_logs_bridge(_mode: &Mode) -> Option { + None +} + +#[cfg(all(feature = "otel", feature = "enterprise"))] +pub(crate) fn init_otlp_tracer(_mode: &Mode) -> Option { + None +} + +pub(crate) fn init_meter_provider(_mode: &Mode) -> OtelProvider { + None +} + +pub fn add_root_flow_job_to_otlp(_queued_job: &QueuedJob, _success: bool) {} diff --git a/backend/windmill-common/src/tracing_init.rs b/backend/windmill-common/src/tracing_init.rs index 13b345e36e..af34540add 100644 --- a/backend/windmill-common/src/tracing_init.rs +++ b/backend/windmill-common/src/tracing_init.rs @@ -7,13 +7,23 @@ */ use const_format::concatcp; + +use std::{ + collections::HashMap, + sync::{Arc, RwLock}, +}; +use tracing::Event; use tracing_appender::non_blocking::{NonBlockingBuilder, WorkerGuard}; +use tracing_subscriber::layer::Context; use tracing_subscriber::{ + filter::Targets, fmt::{format, Layer}, prelude::*, EnvFilter, }; +use crate::utils::Mode; + fn json_layer() -> Layer> { tracing_subscriber::fmt::layer() .json() @@ -34,7 +44,10 @@ pub const LOGS_SERVICE: &str = "logs/services/"; pub const TMP_WINDMILL_LOGS_SERVICE: &str = concatcp!("/tmp/windmill/", LOGS_SERVICE); -pub fn initialize_tracing(hostname: &str) -> WorkerGuard { +pub fn initialize_tracing( + hostname: &str, + mode: &Mode, +) -> (WorkerGuard, crate::otel_ee::OtelProvider) { let style = std::env::var("RUST_LOG_STYLE").unwrap_or_else(|_| "auto".into()); if std::env::var("RUST_LOG").is_ok_and(|x| x == "debug" || x == "info") { @@ -44,7 +57,17 @@ pub fn initialize_tracing(hostname: &str) -> WorkerGuard { ) } - let env_filter = EnvFilter::from_default_env(); + let meter_provider = crate::otel_ee::init_meter_provider(mode); + + #[cfg(all(feature = "otel", feature = "enterprise"))] + let opentelemetry = crate::otel_ee::init_otlp_tracer(mode) + .map(|x| tracing_opentelemetry::layer().with_tracer(x)); + + #[cfg(not(all(feature = "otel", feature = "enterprise")))] + let opentelemetry: Option = None; + + let logs_bridge = crate::otel_ee::init_logs_bridge(&mode); + use tracing_appender::rolling::{RollingFileAppender, Rotation}; let log_dir = format!("{}/{}/", TMP_WINDMILL_LOGS_SERVICE, hostname); @@ -61,39 +84,59 @@ pub fn initialize_tracing(hostname: &str) -> WorkerGuard { .finish(file_appender); let stdout_and_log_file_writer = std::io::stdout.and(log_file_writer); - let ts_base = tracing_subscriber::registry().with(env_filter); + // let job_logs_filter = tracing_subscriber::filter::Targets::new() + // .with_target("windmill:job_log", tracing::Level::TRACE); - #[cfg(feature = "loki")] - let ts_base = { - let (layer, task) = tracing_loki::builder() - .build_url(reqwest::Url::parse("http://127.0.0.1:3100").unwrap()) - .expect("build loki url"); - tokio::spawn(task); - ts_base.with(layer) - }; + let env_filter = EnvFilter::builder() + .with_default_directive(tracing::level_filters::LevelFilter::ERROR.into()) + .from_env_lossy(); + + let ts_base = tracing_subscriber::registry().with(env_filter); match *JSON_FMT { true => ts_base + .with(logs_bridge) + .with(opentelemetry) + // .with(env_filter2.add_directive("windmill:job_log=off".parse().unwrap())) .with( json_layer() .with_writer(stdout_and_log_file_writer) - .flatten_event(true), + .flatten_event(true) + .with_filter( + Targets::new() + .with_target( + "windmill:job_log", + tracing::level_filters::LevelFilter::OFF, + ) + .with_default(tracing::level_filters::LevelFilter::INFO), + ), ) .with(CountingLayer::new()) .init(), false => ts_base + .with(logs_bridge) + .with(opentelemetry) + // .with(env_filter2.add_directive("windmill:job_log=off".parse().unwrap())) .with( compact_layer() .with_writer(stdout_and_log_file_writer) .with_ansi(style.to_lowercase() != "never") .with_file(true) .with_line_number(true) - .with_target(false), + .with_target(false) + .with_filter( + Targets::new() + .with_target( + "windmill:job_log", + tracing::level_filters::LevelFilter::OFF, + ) + .with_default(tracing::level_filters::LevelFilter::INFO), + ), ) .with(CountingLayer::new()) .init(), } - _guard + (_guard, meter_provider) } #[cfg(feature = "flamegraph")] @@ -112,13 +155,6 @@ pub fn setup_flamegraph() -> impl Drop { _guard } -use std::{ - collections::HashMap, - sync::{Arc, RwLock}, -}; -use tracing::Event; -use tracing_subscriber::layer::Context; - lazy_static::lazy_static! { pub static ref LOG_COUNTING_BY_MIN: Arc>> = Arc::new(RwLock::new(HashMap::new())); } @@ -144,22 +180,6 @@ impl CountingLayer { } } -// impl CountingLayer { -// pub fn new() -> Self { -// CountingLayer { counter: Arc::new(Mutex::new(LogCounter::new())) } -// } - -// pub fn get_counts(&self) -> (usize, usize) { -// let counter = self.counter.lock().unwrap(); -// (counter.non_error_count, counter.error_count) -// } - -// pub fn reset_counts(&self) { -// let mut counter = self.counter.lock().unwrap(); -// counter.reset(); -// } -// } - pub const LOG_TIMESTAMP_FMT: &str = "%Y-%m-%d-%H-%M"; impl tracing_subscriber::Layer for CountingLayer diff --git a/backend/windmill-queue/Cargo.toml b/backend/windmill-queue/Cargo.toml index 7ac744a885..8f94de918a 100644 --- a/backend/windmill-queue/Cargo.toml +++ b/backend/windmill-queue/Cargo.toml @@ -43,4 +43,5 @@ bigdecimal.workspace = true axum.workspace = true serde_urlencoded.workspace = true regex.workspace = true -backon.workspace = true \ No newline at end of file +backon.workspace = true +opentelemetry.workspace = true diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index c26a92872f..ef37834400 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -32,7 +32,6 @@ use sqlx::{types::Json, FromRow, Pool, Postgres, Transaction}; #[cfg(feature = "benchmark")] use std::time::Instant; use tokio::{sync::RwLock, time::sleep}; -use tracing::{instrument, Instrument}; use ulid::Ulid; use uuid::Uuid; use windmill_audit::audit_ee::{audit_log, AuditAuthor}; @@ -454,7 +453,6 @@ where } } -#[instrument(level = "trace", skip_all)] pub async fn add_completed_job_error( db: &Pool, queued_job: &QueuedJob, @@ -510,7 +508,6 @@ lazy_static::lazy_static! { pub static ref GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE: Option = std::env::var("GLOBAL_ERROR_HANDLER_PATH_IN_ADMINS_WORKSPACE").ok(); } -#[instrument(level = "trace", skip_all, name = "add_completed_job")] pub async fn add_completed_job( db: &Pool, queued_job: &QueuedJob, @@ -643,9 +640,7 @@ pub async fn add_completed_job( .fetch_one(&mut *tx) .await .map_err(|e| Error::InternalErr(format!("Could not add completed job {job_id}: {e:#}")))?; - // tracing::error!("2 {:?}", start.elapsed()); - // add_time!(bench, "add_completed_job query END"); if !queued_job.is_flow_step { if _duration > 500 @@ -1259,7 +1254,6 @@ pub async fn send_error_to_workspace_handler<'a, 'c, T: Serialize + Send + Sync> Ok(()) } -#[instrument(level = "trace", skip_all)] pub async fn handle_maybe_scheduled_job<'c>( db: &Pool, job: &QueuedJob, @@ -2429,7 +2423,6 @@ async fn extract_result_from_job_result( } } -#[instrument(level = "trace", skip_all)] pub async fn delete_job<'c>( mut tx: Transaction<'c, Postgres>, w_id: &str, @@ -4039,7 +4032,6 @@ pub async fn push<'c, 'd>( script_path.as_ref().map(|x| x.as_str()), Some(hm), ) - .instrument(tracing::info_span!("job_run", email = &email)) .await?; } diff --git a/backend/windmill-worker/Cargo.toml b/backend/windmill-worker/Cargo.toml index cb64757a7f..22640e577a 100644 --- a/backend/windmill-worker/Cargo.toml +++ b/backend/windmill-worker/Cargo.toml @@ -19,6 +19,7 @@ flow_testing = [] cloud = [] sqlx = [] deno_core = ["dep:deno_fetch", "dep:deno_webidl", "dep:deno_web", "dep:deno_net", "dep:deno_console", "dep:deno_url", "dep:deno_core", "dep:deno_ast", "dep:deno_tls"] +otel = ["windmill-common/otel", "dep:opentelemetry"] [dependencies] windmill-queue.workspace = true @@ -92,6 +93,9 @@ yaml-rust.workspace = true swc_ecma_parser.workspace = true backon.workspace = true +opentelemetry = { workspace = true, optional = true } + + [build-dependencies] deno_fetch = { workspace = true, optional = true } deno_webidl = { workspace = true, optional = true } diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 31d8ecc4dd..f36954819d 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -866,11 +866,11 @@ fn tentatively_improve_error(err: Error, executable: &str) -> Error { #[cfg(windows)] let err_msg = "program not found"; - if err - .to_string() - .contains(&err_msg) - { - return Error::InternalErr(format!("Executable {executable} not found on worker. PATH: {}", *PATH_ENV)); + if err.to_string().contains(&err_msg) { + return Error::InternalErr(format!( + "Executable {executable} not found on worker. PATH: {}", + *PATH_ENV + )); } return err; } diff --git a/backend/windmill-worker/src/handle_child.rs b/backend/windmill-worker/src/handle_child.rs index 4d09c36abc..cb35070a40 100644 --- a/backend/windmill-worker/src/handle_child.rs +++ b/backend/windmill-worker/src/handle_child.rs @@ -87,7 +87,7 @@ async fn kill_process_tree(pid: Option) -> Result<(), String> { /// - update the `last_line` and `logs` strings with the program output /// - update "queue"."last_ping" every five seconds /// - kill process if we exceed timeout or "queue"."canceled" is set -#[tracing::instrument(level = "trace", skip_all)] +#[tracing::instrument(name="run_subprocess", level = "info", skip_all, fields(otel.name = %child_name))] pub async fn handle_child( job_id: &Uuid, db: &Pool, @@ -435,8 +435,6 @@ pub async fn handle_child( } } - - async fn get_mem_peak(pid: Option, nsjail: bool) -> i32 { if pid.is_none() { return -1; @@ -692,7 +690,10 @@ fn child_joined_output_stream( let stdout = BufReader::new(stdout).lines(); let stderr = BufReader::new(stderr).lines(); - stream::select(lines_to_stream(stderr, true), lines_to_stream(stdout, false)) + stream::select( + lines_to_stream(stderr, true), + lines_to_stream(stdout, false), + ) } pub fn lines_to_stream( @@ -702,14 +703,10 @@ pub fn lines_to_stream( stream::poll_fn(move |cx| { std::pin::Pin::new(&mut lines) .poll_next_line(cx) - .map(|result| { - process_streaming_log_lines(result, stderr) - }) + .map(|result| process_streaming_log_lines(result, stderr)) }) } - - pub fn process_status(status: ExitStatus) -> error::Result<()> { if status.success() { Ok(()) diff --git a/backend/windmill-worker/src/result_processor.rs b/backend/windmill-worker/src/result_processor.rs index 25bfade2e8..a897459224 100644 --- a/backend/windmill-worker/src/result_processor.rs +++ b/backend/windmill-worker/src/result_processor.rs @@ -1,3 +1,6 @@ +#[cfg(feature = "otel")] +use opentelemetry::trace::FutureExt; + use serde::Serialize; use sqlx::{types::Json, Pool, Postgres}; use std::{ @@ -7,6 +10,9 @@ use std::{ Arc, }, }; +use tracing::{field, Instrument}; +#[cfg(not(feature = "otel"))] +use windmill_common::otel_ee::FutureExt; use uuid::Uuid; @@ -85,7 +91,49 @@ pub fn start_background_processor( JobKind::Dependencies | JobKind::FlowDependencies ); - handle_receive_completed_job( + let success = jc.success; + + let span = tracing::span!( + tracing::Level::INFO, + "job_postprocessing", + job_id = %jc.job.id, root_job = field::Empty, workspace_id = %jc.job.workspace_id, worker = %worker_name,tag = %jc.job.tag, + // hostname = %hostname, + language = field::Empty, + script_path = field::Empty, + flow_step_id = field::Empty, + parent_job = field::Empty, + otel.name = field::Empty + ); + let rj = if let Some(root_job) = jc.job.root_job { + root_job + } else { + jc.job.id + }; + windmill_common::otel_ee::set_span_parent(&span, &rj); + + if let Some(lg) = jc.job.language.as_ref() { + span.record("language", lg.as_str()); + } + if let Some(step_id) = jc.job.flow_step_id.as_ref() { + span.record( + "otel.name", + format!("job_postprocessing {}", step_id).as_str(), + ); + span.record("flow_step_id", step_id.as_str()); + } else { + span.record("otel.name", "job postprocessing"); + } + if let Some(parent_job) = jc.job.parent_job.as_ref() { + span.record("parent_job", parent_job.to_string().as_str()); + } + if let Some(script_path) = jc.job.script_path.as_ref() { + span.record("script_path", script_path.as_str()); + } + if let Some(root_job) = jc.job.root_job.as_ref() { + span.record("root_job", root_job.to_string().as_str()); + } + + let root_job = handle_receive_completed_job( jc, &base_internal_url, &db, @@ -96,8 +144,14 @@ pub fn start_background_processor( #[cfg(feature = "benchmark")] &mut bench, ) + .instrument(span) .await; + if let Some(root_job) = root_job { + windmill_common::otel_ee::add_root_flow_job_to_otlp(&root_job, success); + tracing::error!(job_id = %root_job.id, parent_job = ?root_job.parent_job, "ADDDED root job completed"); + } + if is_init_script_and_failure { tracing::error!("init script errored, exiting"); killpill_tx.send(()).unwrap_or_default(); @@ -199,7 +253,11 @@ async fn send_job_completed( token, duration, }; - job_completed_tx.send(jc).await.expect("send job completed") + job_completed_tx + .send(jc) + .with_context(windmill_common::otel_ee::otel_ctx()) + .await + .expect("send job completed") } pub async fn process_result( @@ -271,6 +329,7 @@ pub async fn process_result( token, duration, ) + .with_context(windmill_common::otel_ee::otel_ctx()) .await; Ok(true) } @@ -315,6 +374,7 @@ pub async fn process_result( token, duration, ) + .with_context(windmill_common::otel_ee::otel_ctx()) .await; Ok(false) } @@ -330,7 +390,7 @@ pub async fn handle_receive_completed_job( worker_name: &str, job_completed_tx: Sender, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, -) { +) -> Option> { let token = jc.token.clone(); let workspace = jc.job.workspace_id.clone(); let client = AuthedClient { @@ -342,7 +402,7 @@ pub async fn handle_receive_completed_job( let job = jc.job.clone(); let mem_peak = jc.mem_peak.clone(); let canceled_by = jc.canceled_by.clone(); - if let Err(err) = process_completed_job( + match process_completed_job( jc, &client, db, @@ -355,26 +415,29 @@ pub async fn handle_receive_completed_job( ) .await { - handle_job_error( - db, - &client, - job.as_ref(), - mem_peak, - canceled_by, - err, - false, - same_worker_tx.clone(), - &worker_dir, - worker_name, - job_completed_tx, - #[cfg(feature = "benchmark")] - bench, - ) - .await; + Err(err) => { + handle_job_error( + db, + &client, + job.as_ref(), + mem_peak, + canceled_by, + err, + false, + same_worker_tx.clone(), + &worker_dir, + worker_name, + job_completed_tx, + #[cfg(feature = "benchmark")] + bench, + ) + .await; + None + } + Ok(r) => r, } } -#[tracing::instrument(name = "completed_job", level = "info", skip_all, fields(job_id = %job.id))] pub async fn process_completed_job( JobCompleted { job, result, mem_peak, success, cached_res_path, canceled_by, duration, .. }: JobCompleted, client: &AuthedClient, @@ -384,7 +447,7 @@ pub async fn process_completed_job( worker_name: &str, job_completed_tx: Sender, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, -) -> windmill_common::error::Result<()> { +) -> windmill_common::error::Result>> { if success { // println!("bef completed job{:?}", SystemTime::now()); if let Some(cached_path) = cached_res_path { @@ -414,8 +477,8 @@ pub async fn process_completed_job( if is_flow_step { if let Some(parent_job) = parent_job { - tracing::info!(parent_flow = %parent_job, subflow = %job_id, "updating flow status (2)"); - update_flow_status_after_job_completion( + // tracing::info!(parent_flow = %parent_job, subflow = %job_id, "updating flow status (2)"); + let r = update_flow_status_after_job_completion( db, client, parent_job, @@ -434,9 +497,10 @@ pub async fn process_completed_job( ) .warn_after_seconds(10) .await?; + add_time!(bench, "updated flow status END"); + return Ok(r); } } - add_time!(bench, "updated flow status END"); } else { let result = add_completed_job_error( db, @@ -454,7 +518,7 @@ pub async fn process_completed_job( if job.is_flow_step { if let Some(parent_job) = job.parent_job { tracing::error!(parent_flow = %parent_job, subflow = %job.id, "process completed job error, updating flow status"); - update_flow_status_after_job_completion( + let r = update_flow_status_after_job_completion( db, client, parent_job, @@ -473,10 +537,11 @@ pub async fn process_completed_job( ) .warn_after_seconds(10) .await?; + return Ok(r); } } } - Ok(()) + return Ok(None); } #[tracing::instrument(name = "job_error", level = "info", skip_all, fields(job_id = %job.id))] diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 238fd6a951..f179396123 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -6,6 +6,10 @@ * LICENSE-AGPL for a copy of the license. */ +// #[cfg(feature = "otel")] +// use opentelemetry::{global, KeyValue}; + + use windmill_common::{ apps::AppScriptId, auth::{fetch_authed_from_permissioned_as, JWTAuthClaims, JobPerms, JWT_SECRET}, @@ -25,7 +29,7 @@ use const_format::concatcp; #[cfg(feature = "prometheus")] use prometheus::IntCounter; -use tracing::Instrument; +use tracing::{field, Instrument}; #[cfg(feature = "prometheus")] use windmill_common::METRICS_DEBUG_ENABLED; #[cfg(feature = "prometheus")] @@ -88,30 +92,12 @@ use tokio::{ use rand::Rng; use crate::{ - ansible_executor::handle_ansible_job, - bash_executor::{handle_bash_job, handle_powershell_job}, - bun_executor::handle_bun_job, - common::{ + ansible_executor::handle_ansible_job, bash_executor::{handle_bash_job, handle_powershell_job}, bun_executor::handle_bun_job, common::{ build_args_map, get_cached_resource_value_if_valid, get_reserved_variables, hash_args, update_worker_ping_for_failed_init_script, OccupancyMetrics, - }, - deno_executor::handle_deno_job, - go_executor::handle_go_job, - graphql_executor::do_graphql, - handle_child::SLOW_LOGS, - handle_job_error, - job_logger::NO_LOGS_AT_ALL, - js_eval::{eval_fetch_timeout, transpile_ts}, - mysql_executor::do_mysql, - pg_executor::do_postgresql, - php_executor::handle_php_job, - python_executor::handle_python_job, - result_processor::{process_result, start_background_processor}, - rust_executor::handle_rust_job, - worker_flow::{handle_flow, update_flow_status_in_progress, Step}, - worker_lockfiles::{ + }, deno_executor::handle_deno_job, go_executor::handle_go_job, graphql_executor::do_graphql, handle_child::SLOW_LOGS, handle_job_error, job_logger::NO_LOGS_AT_ALL, js_eval::{eval_fetch_timeout, transpile_ts}, mysql_executor::do_mysql, pg_executor::do_postgresql, php_executor::handle_php_job, python_executor::handle_python_job, result_processor::{process_result, start_background_processor}, rust_executor::handle_rust_job, worker_flow::{handle_flow, update_flow_status_in_progress, Step}, worker_lockfiles::{ handle_app_dependency_job, handle_dependency_job, handle_flow_dependency_job, - }, + } }; use backon::ConstantBuilder; @@ -720,7 +706,10 @@ fn add_outstanding_wait_time( }.in_current_span()); } -#[tracing::instrument(name = "worker", level = "info", skip_all, fields(worker = %worker_name, hostname = %hostname))] +// struct WorkerMtrics { +// job_ +// } + pub async fn run_worker( db: &Pool, hostname: &str, @@ -736,6 +725,7 @@ pub async fn run_worker( #[cfg(not(feature = "enterprise"))] if !*DISABLE_NSJAIL { tracing::warn!( + worker = %worker_name, hostname = %hostname, "NSJAIL to sandbox process in untrusted environments is an enterprise feature but allowed to be used for testing purposes" ); } @@ -743,10 +733,10 @@ pub async fn run_worker( let start_time = Instant::now(); let worker_dir = format!("{TMP_DIR}/{worker_name}"); - tracing::debug!(worker_dir = %worker_dir, "Creating worker dir"); + tracing::debug!(worker = %worker_name, hostname = %hostname, worker_dir = %worker_dir, "Creating worker dir"); if let Some(ref netrc) = *NETRC { - tracing::info!("Writing netrc at {}/.netrc", HOME_ENV.as_str()); + tracing::info!(worker = %worker_name, hostname = %hostname, "Writing netrc at {}/.netrc", HOME_ENV.as_str()); write_file(&HOME_ENV, ".netrc", netrc).expect("could not write netrc"); } @@ -961,6 +951,16 @@ pub async fn run_worker( None }; + + // let worker_resource = &[ + // KeyValue::new("hostname", hostname.to_string()), + // KeyValue::new("worker", worker_name.to_string()), + // ]; + // // Create a meter from the above MeterProvider. + // let meter = global::meter("windmill"); + // let counter = meter.u64_counter("jobs.execution").build(); + + let mut occupancy_metrics = OccupancyMetrics::new(start_time); let mut jobs_executed = 0; @@ -1016,6 +1016,7 @@ pub async fn run_worker( IS_READY.store(true, Ordering::Relaxed); tracing::info!( + worker = %worker_name, hostname = %hostname, "listening for jobs, WORKER_GROUP: {}, config: {:?}", *WORKER_GROUP, WORKER_CONFIG.read().await @@ -1051,7 +1052,7 @@ pub async fn run_worker( if i_worker == 1 { if let Err(e) = queue_init_bash_maybe(db, same_worker_tx.clone(), &worker_name).await { killpill_tx.send(()).unwrap_or_default(); - tracing::error!("Error queuing init bash script for worker {worker_name}: {e:#}"); + tracing::error!(worker = %worker_name, hostname = %hostname, "Error queuing init bash script for worker {worker_name}: {e:#}"); return; } } @@ -1084,7 +1085,7 @@ pub async fn run_worker( #[cfg(feature = "enterprise")] { if let Ok(_) = killpill_rx.try_recv() { - tracing::info!("killpill received on worker waiting for valid key"); + tracing::info!(worker = %worker_name, hostname = %hostname, "killpill received on worker waiting for valid key"); job_completed_tx .0 .send(SendResult::Kill) @@ -1096,6 +1097,7 @@ pub async fn run_worker( if !valid_key { tracing::error!( + worker = %worker_name, hostname = %hostname, "Invalid license key, workers require a valid license key, sleeping for 30s waiting for valid key to be set" ); tokio::time::sleep(Duration::from_secs(10)).await; @@ -1109,7 +1111,7 @@ pub async fn run_worker( #[cfg(feature = "prometheus")] if let Some(wk) = worker_busy.as_ref() { wk.set(0); - tracing::debug!("set worker busy to 0"); + tracing::debug!(worker = %worker_name, hostname = %hostname, "set worker busy to 0"); } occupancy_metrics.running_job_started_at = None; @@ -1121,7 +1123,7 @@ pub async fn run_worker( .try_into() .unwrap(), ); - tracing::debug!("set uptime metric"); + tracing::debug!(worker = %worker_name, hostname = %hostname, "set uptime metric"); } if last_ping.elapsed().as_secs() > NUM_SECS_PING { @@ -1165,15 +1167,19 @@ pub async fn run_worker( ) .notify(|err, dur| { tracing::error!( + worker = %worker_name, hostname = %hostname, "retrying updating worker ping in {dur:#?}, err: {err:#?}" ); }) .sleep(tokio::time::sleep) .await { - tracing::error!("failed to update worker ping, exiting: {}", e); + tracing::error!( + worker = %worker_name, hostname = %hostname, + "failed to update worker ping, exiting: {}", e); killpill_tx.send(()).unwrap_or_default(); } tracing::info!( + worker = %worker_name, hostname = %hostname, "ping update, memory: container={}MB, windmill={}MB", memory_usage.unwrap_or_default() / (1024 * 1024), wm_memory_usage.unwrap_or_default() / (1024 * 1024) @@ -1185,16 +1191,18 @@ pub async fn run_worker( if (jobs_executed as u32 + vacuum_shift) % VACUUM_PERIOD == 0 { let db2 = db.clone(); let current_span = tracing::Span::current(); + let worker_name = worker_name.clone(); + let hostname = hostname.to_string(); tokio::task::spawn( (async move { - tracing::info!("vacuuming queue and completed_job"); + tracing::info!(worker = %worker_name, hostname = %hostname, "vacuuming queue"); if let Err(e) = sqlx::query!("VACUUM (skip_locked) queue") .execute(&db2) .await { - tracing::error!("failed to vacuum queue: {}", e); + tracing::error!(worker = %worker_name, hostname = %hostname, "failed to vacuum queue: {}", e); } - tracing::info!("vacuumed queue and completed_job"); + tracing::info!(worker = %worker_name, hostname = %hostname, "vacuumed queue"); }) .instrument(current_span), ); @@ -1230,6 +1238,7 @@ pub async fn run_worker( if let Ok(same_worker_job) = same_worker_rx.try_recv() { same_worker_queue_size.fetch_sub(1, Ordering::SeqCst); tracing::debug!( + worker = %worker_name, hostname = %hostname, "received {} from same worker channel", same_worker_job.job_id ); @@ -1242,6 +1251,7 @@ pub async fn run_worker( .map_err(|_| Error::InternalErr("Impossible to fetch same_worker job".to_string())); if r.is_err() && !same_worker_job.recoverable { tracing::error!( + worker = %worker_name, hostname = %hostname, "failed to fetch same_worker job on a non recoverable job, exiting" ); job_completed_tx @@ -1255,7 +1265,7 @@ pub async fn run_worker( } } else if let Ok(_) = killpill_rx.try_recv() { if !killed_but_draining_same_worker_jobs { - tracing::info!("received killpill for worker {}, jobs are not pulled anymore except same_worker jobs", i_worker); + tracing::info!(worker = %worker_name, hostname = %hostname, "received killpill for worker {}, jobs are not pulled anymore except same_worker jobs", i_worker); killed_but_draining_same_worker_jobs = true; job_completed_tx .0 @@ -1266,10 +1276,10 @@ pub async fn run_worker( continue; } else if killed_but_draining_same_worker_jobs { if job_completed_processor_is_done.load(Ordering::SeqCst) { - tracing::info!("all running jobs have completed and all completed jobs have been fully processed, exiting"); + tracing::info!(worker = %worker_name, hostname = %hostname, "all running jobs have completed and all completed jobs have been fully processed, exiting"); break; } else { - tracing::info!("there may be same_worker jobs to process later, waiting for job_completed_processor to finish progressing all remaining flows before exiting"); + tracing::info!(worker = %worker_name, hostname = %hostname, "there may be same_worker jobs to process later, waiting for job_completed_processor to finish progressing all remaining flows before exiting"); tokio::time::sleep(Duration::from_millis(200)).await; continue; } @@ -1294,7 +1304,7 @@ pub async fn run_worker( if !agent_mode && duration_pull_s > 0.5 { let empty = job.as_ref().is_ok_and(|x| x.0.is_none()); - tracing::warn!("pull took more than 0.5s ({duration_pull_s}), this is a sign that the database is VERY undersized for this load. empty: {empty}, err: {err_pull}"); + tracing::warn!(worker = %worker_name, hostname = %hostname, "pull took more than 0.5s ({duration_pull_s}), this is a sign that the database is VERY undersized for this load. empty: {empty}, err: {err_pull}"); #[cfg(feature = "prometheus")] if empty { if let Some(wp) = worker_pull_over_500_counter_empty.as_ref() { @@ -1305,7 +1315,7 @@ pub async fn run_worker( } } else if !agent_mode && duration_pull_s > 0.1 { let empty = job.as_ref().is_ok_and(|x| x.0.is_none()); - tracing::warn!("pull took more than 0.1s ({duration_pull_s}) this is a sign that the database is undersized for this load. empty: {empty}, err: {err_pull}"); + tracing::warn!(worker = %worker_name, hostname = %hostname, "pull took more than 0.1s ({duration_pull_s}) this is a sign that the database is undersized for this load. empty: {empty}, err: {err_pull}"); #[cfg(feature = "prometheus")] if empty { if let Some(wp) = worker_pull_over_100_counter_empty.as_ref() { @@ -1359,7 +1369,7 @@ pub async fn run_worker( last_executed_job = None; jobs_executed += 1; - tracing::debug!("started handling of job {}", job.id); + tracing::debug!(worker = %worker_name, hostname = %hostname, "started handling of job {}", job.id); if matches!(job.job_kind, JobKind::Script | JobKind::Preview) { if !dedicated_workers.is_empty() { @@ -1424,6 +1434,11 @@ pub async fn run_worker( ) .await; + // counter.add( + // 1, + // worker_resource + // ); + #[cfg(feature = "prometheus")] let _timer = register_metric( &WORKER_EXECUTION_DURATION, @@ -1511,6 +1526,40 @@ pub async fn run_worker( let PulledJob { job, raw_code, raw_lock, raw_flow } = job; let arc_job = Arc::new(job); add_time!(bench, "handle_queued_job START"); + + + let span = tracing::span!(tracing::Level::INFO, "job", + job_id = %arc_job.id, root_job = field::Empty, workspace_id = %arc_job.workspace_id, worker = %worker_name, hostname = %hostname, tag = %arc_job.tag, + language = field::Empty, + script_path = field::Empty, flow_step_id = field::Empty, parent_job = field::Empty, + otel.name = field::Empty); + let rj = if let Some(root_job) = arc_job.root_job { + root_job + } else { + arc_job.id + }; + if let Some(lg) = arc_job.language.as_ref() { + span.record("language", lg.as_str()); + } + if let Some(step_id) = arc_job.flow_step_id.as_ref() { + span.record("otel.name", format!("job {}", step_id).as_str()); + span.record("flow_step_id", step_id.as_str()); + } else { + span.record("otel.name", "job"); + } + if let Some(parent_job) = arc_job.parent_job.as_ref() { + span.record("parent_job", parent_job.to_string().as_str()); + } + if let Some(script_path) = arc_job.script_path.as_ref() { + span.record("script_path", script_path.as_str()); + } + if let Some(root_job) = arc_job.root_job.as_ref() { + span.record("root_job", root_job.to_string().as_str()); + } + + windmill_common::otel_ee::set_span_parent(&span, &rj); + // span.context().span().add_event_with_timestamp("job created".to_string(), arc_job.created_at.into(), vec![]); + match handle_queued_job( arc_job.clone(), raw_code, @@ -1529,6 +1578,7 @@ pub async fn run_worker( #[cfg(feature = "benchmark")] &mut bench, ) + .instrument(span) .await { Err(err) => { @@ -1568,6 +1618,8 @@ pub async fn run_worker( _ => {} } + + #[cfg(feature = "prometheus")] if let Some(duration) = _timer.map(|x| x.stop_and_record()) { register_metric( @@ -1607,7 +1659,7 @@ pub async fn run_worker( if let Some(secs) = *EXIT_AFTER_NO_JOB_FOR_SECS { if let Some(lj) = last_executed_job { if lj.elapsed().as_secs() > secs { - tracing::info!("no job for {} seconds, exiting", secs); + tracing::info!(worker = %worker_name, hostname = %hostname, "no job for {} seconds, exiting", secs); break; } } else { @@ -1638,12 +1690,12 @@ pub async fn run_worker( }); } Err(err) => { - tracing::error!("Failed to pull jobs: {}", err); + tracing::error!(worker = %worker_name, hostname = %hostname, "Failed to pull jobs: {}", err); } }; } - tracing::info!("worker {} exiting", worker_name); + tracing::info!(worker = %worker_name, hostname = %hostname, "worker {} exiting", worker_name); #[cfg(feature = "benchmark")] { @@ -1658,22 +1710,24 @@ pub async fn run_worker( if has_dedicated_workers { for handle in dedicated_handles { if let Err(e) = handle.await { - tracing::error!("error in dedicated worker waiting for it to end: {:?}", e) + tracing::error!(worker = %worker_name, hostname = %hostname, "error in dedicated worker waiting for it to end: {:?}", e) } } - tracing::info!("all dedicated workers have exited"); + tracing::info!(worker = %worker_name, hostname = %hostname, "all dedicated workers have exited"); } drop(job_completed_tx); - tracing::info!("waiting for job_completed_processor to finish processing remaining jobs"); + tracing::info!(worker = %worker_name, hostname = %hostname, "waiting for job_completed_processor to finish processing remaining jobs"); if let Err(e) = send_result.await { tracing::error!("error in awaiting send_result process: {e:?}") } - tracing::info!("worker {} exited", worker_name); - tracing::info!("number of jobs executed: {}", jobs_executed); + tracing::info!(worker = %worker_name, hostname = %hostname, "worker {} exited", worker_name); + tracing::info!(worker = %worker_name, hostname = %hostname, "number of jobs executed: {}", jobs_executed); } + + async fn queue_init_bash_maybe<'c>( db: &Pool, same_worker_tx: SameWorkerSender, @@ -1798,7 +1852,6 @@ pub struct PreviousResult<'a> { pub previous_result: Option<&'a RawValue>, } -#[tracing::instrument(name = "job", level = "info", skip_all, fields(job_id = %job.id))] async fn handle_queued_job( job: Arc, raw_code: Option, @@ -1816,6 +1869,9 @@ async fn handle_queued_job( occupancy_metrics: &mut OccupancyMetrics, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, ) -> windmill_common::error::Result { + // Extract the active span from the context + + if job.canceled { return Err(Error::JsonErr(canceled_job_to_result(&job))); } @@ -2159,6 +2215,7 @@ async fn handle_queued_job( } } + pub fn build_envs( envs: Option>, ) -> windmill_common::error::Result> { diff --git a/backend/windmill-worker/src/worker_flow.rs b/backend/windmill-worker/src/worker_flow.rs index 74949d0bcf..b5770166c5 100644 --- a/backend/windmill-worker/src/worker_flow.rs +++ b/backend/windmill-worker/src/worker_flow.rs @@ -77,45 +77,34 @@ pub async fn update_flow_status_after_job_completion( worker_name: &str, job_completed_tx: Sender, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, -) -> error::Result<()> { +) -> error::Result>> { // this is manual tailrecursion because async_recursion blows up the stack - // todo!(); potentially_crash_for_testing(); - let mut rec = update_flow_status_after_job_completion_internal( - db, - client, + let mut rec = RecUpdateFlowStatusAfterJobCompletion { flow, - job_id_for_status, - w_id, + job_id_for_status: job_id_for_status.clone(), success, result, - unrecoverable, - same_worker_tx.clone(), - worker_dir, stop_early_override, - false, - worker_name, - job_completed_tx.clone(), - #[cfg(feature = "benchmark")] - bench, - ) - .await?; - while let Some(nrec) = rec { + skip_error_handler: false, + }; + let mut unrecoverable = unrecoverable; + loop { potentially_crash_for_testing(); - rec = match update_flow_status_after_job_completion_internal( + let nrec = match update_flow_status_after_job_completion_internal( db, client, - nrec.flow, - &nrec.job_id_for_status, + rec.flow, + &rec.job_id_for_status, w_id, - nrec.success, - nrec.result, - false, + rec.success, + rec.result, + unrecoverable, same_worker_tx.clone(), worker_dir, - nrec.stop_early_override, - nrec.skip_error_handler, + rec.stop_early_override, + rec.skip_error_handler, worker_name, job_completed_tx.clone(), #[cfg(feature = "benchmark")] @@ -125,12 +114,12 @@ pub async fn update_flow_status_after_job_completion( { Ok(j) => j, Err(e) => { - tracing::error!("Error while updating flow status of {} after completion of {}, updating flow status again with error: {e:#}", nrec.flow,&nrec.job_id_for_status); + tracing::error!("Error while updating flow status of {} after completion of {}, updating flow status again with error: {e:#}", rec.flow, &rec.job_id_for_status); update_flow_status_after_job_completion_internal( db, client, - nrec.flow, - &nrec.job_id_for_status, + rec.flow, + &rec.job_id_for_status, w_id, false, Arc::new(to_raw_value(&Json(&WrappedError { @@ -139,8 +128,8 @@ pub async fn update_flow_status_after_job_completion( true, same_worker_tx.clone(), worker_dir, - nrec.stop_early_override, - nrec.skip_error_handler, + rec.stop_early_override, + rec.skip_error_handler, worker_name, job_completed_tx.clone(), #[cfg(feature = "benchmark")] @@ -148,9 +137,33 @@ pub async fn update_flow_status_after_job_completion( ) .await? } + }; + unrecoverable = false; + match nrec { + UpdateFlowStatusAfterJobCompletion::Done(job) => { + add_time!(bench, "update flow status internal END"); + return Ok(Some(job)); + } + UpdateFlowStatusAfterJobCompletion::Rec(nrec) => { + rec = nrec; + }, + UpdateFlowStatusAfterJobCompletion::NonLastParallelBranch => { + add_time!(bench, "update flow status internal END"); + return Ok(None); + }, + UpdateFlowStatusAfterJobCompletion::NotDone => { + add_time!(bench, "update flow status internal END"); + return Ok(None); + } } } - Ok(()) +} + +pub enum UpdateFlowStatusAfterJobCompletion { + Rec(RecUpdateFlowStatusAfterJobCompletion), + Done(Arc), + NotDone, + NonLastParallelBranch, } pub struct RecUpdateFlowStatusAfterJobCompletion { flow: uuid::Uuid, @@ -188,7 +201,7 @@ pub async fn update_flow_status_after_job_completion_internal( worker_name: &str, job_completed_tx: Sender, #[cfg(feature = "benchmark")] bench: &mut BenchmarkIter, -) -> error::Result> { +) -> error::Result { add_time!(bench, "update flow status internal START"); let ( should_continue_flow, @@ -603,7 +616,7 @@ pub async fn update_flow_status_after_job_completion_internal( ); } add_time!(bench, "non final parallel flow finished"); - return Ok(None); + return Ok(UpdateFlowStatusAfterJobCompletion::NonLastParallelBranch); } } FlowStatusModule::InProgress { @@ -1148,7 +1161,7 @@ pub async fn update_flow_status_after_job_completion_internal( if let Some(parent_job) = flow_job.parent_job { tracing::info!(subflow_id = %flow_job.id, parent_id = %parent_job, "subflow is finished, updating parent flow status"); - return Ok(Some(RecUpdateFlowStatusAfterJobCompletion { + return Ok(UpdateFlowStatusAfterJobCompletion::Rec(RecUpdateFlowStatusAfterJobCompletion { flow: parent_job, job_id_for_status: flow, success: success && !is_failure_step, @@ -1162,9 +1175,9 @@ pub async fn update_flow_status_after_job_completion_internal( })); } } - Ok(None) + Ok(UpdateFlowStatusAfterJobCompletion::Done(flow_job)) } else { - Ok(None) + Ok(UpdateFlowStatusAfterJobCompletion::NotDone) } } diff --git a/frontend/src/lib/components/AuthSettings.svelte b/frontend/src/lib/components/AuthSettings.svelte new file mode 100644 index 0000000000..56f6d625de --- /dev/null +++ b/frontend/src/lib/components/AuthSettings.svelte @@ -0,0 +1,255 @@ + + +
+ + SSO + OAuth + SCIM/SAML + +
+ +
+ {#if tab === 'sso'} + {#if !$enterpriseLicense || $enterpriseLicense.endsWith('_pro')} + + Without EE, the number of SSO users is limited to 10. SCIM/SAML is available on EE + + {/if} + +
+
+ When at least one of the below options is set, users will be able to login to Windmill via + their third-party account. +
To test SSO, the recommended workflow is to to save the settings and try to login in + an incognito window. + Learn more
+
+
+ + + + + + + + + + + + {#each Object.keys(oauths) as k} + {#if !['authelia', 'authentik', 'google', 'microsoft', 'github', 'gitlab', 'jumpcloud', 'okta', 'keycloak', 'slack', 'kanidm', 'zitadel'].includes(k) && 'login_config' in oauths[k]} + {#if oauths[k]} +
+
+ + + { + delete oauths[k] + oauths = { ...oauths } + }} + /> +
+
+ + + + {#if !windmillBuiltins.includes(k) && k != 'slack'} + + {/if} +
+
+ {/if} + {/if} + {/each} +
+
+ + +
+
+ +
+ {:else if tab === 'oauth'} +
+ When one of the below options is set, you will be able to create a specific resource + containing a token automatically generated by the third-party provider. +
+ To test it after setting an oauth client, go to the Resources menu and create a new one of the + type of your oauth client (i.e. a 'github' resource if you set Github OAuth). +
Learn more
+
+
+ +
+ + {#each Object.keys(oauths) as k} + {#if oauths[k] && !('login_config' in oauths[k])} + {#if !['slack'].includes(k) && oauths[k]} +
+
+ + + { + delete oauths[k] + oauths = { ...oauths } + }} + /> +
+
+ + + {#if !windmillBuiltins.includes(k) && k != 'slack'} + + {/if} + {#if k == 'snowflake_oauth'} + + {/if} +
+
+ {/if} + {/if} + {/each} + +
+ + {#if oauth_name == 'custom'} + + {:else} + + {/if} + +
+ {:else if tab == 'scim'} + + {/if} +
diff --git a/frontend/src/lib/components/InstanceSetting.svelte b/frontend/src/lib/components/InstanceSetting.svelte new file mode 100644 index 0000000000..1bac352d5b --- /dev/null +++ b/frontend/src/lib/components/InstanceSetting.svelte @@ -0,0 +1,783 @@ + + +{#if (!setting.cloudonly || isCloudHosted()) && showSetting(setting.key, $values) && !(setting.hiddenIfNull && $values[setting.key] == null)} + {#if setting.ee_only != undefined && !$enterpriseLicense} +
+ + EE only {#if setting.ee_only != ''}{setting.ee_only}{/if} +
+ {/if} + +