From c6dc19a893cf53c220c43c8eaa2ab7fcc8a128a7 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Sun, 8 Dec 2024 11:05:09 +0100 Subject: [PATCH] improve otel --- backend/ee-repo-ref.txt | 2 +- backend/src/main.rs | 15 +++++++++++++-- backend/src/monitor.rs | 16 ++++++++++------ backend/tests/worker.rs | 1 + backend/windmill-common/src/otel_ee.rs | 5 +++-- backend/windmill-common/src/tracing_init.rs | 7 ++++--- backend/windmill-worker/src/handle_child.rs | 10 ++++++---- backend/windmill-worker/src/job_logger_ee.rs | 1 + 8 files changed, 39 insertions(+), 18 deletions(-) diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index dbe12a2939..32341c9e08 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -1d7ee3745e81082c196c4aea0392c88e3de6f9d4 \ No newline at end of file +a07bc62582c809457f1c945d6cda145770b94d04 \ No newline at end of file diff --git a/backend/src/main.rs b/backend/src/main.rs index 80edff1ac9..e76dfa79f6 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -8,7 +8,7 @@ use anyhow::Context; use monitor::{ - load_otel, reload_delete_logs_periodically_setting, reload_indexer_config, + load_base_url, 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, }; @@ -348,10 +348,21 @@ async fn windmill_main() -> anyhow::Result<()> { let db = windmill_common::connect_db(server_mode, indexer_mode).await?; load_otel(&db).await; + tracing::info!("Database connected"); + let environment = load_base_url(&db) + .await + .unwrap_or_else(|_| "local".to_string()) + .trim_start_matches("https://") + .trim_start_matches("http://") + .split(".") + .next() + .unwrap_or_else(|| "local") + .to_string(); + #[cfg(not(feature = "flamegraph"))] - let _guard = windmill_common::tracing_init::initialize_tracing(&hostname, &mode); + let _guard = windmill_common::tracing_init::initialize_tracing(&hostname, &mode, &environment); let num_version = sqlx::query_scalar!("SELECT version()").fetch_one(&db).await; diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 9acbbb1378..c11d60621e 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -1397,7 +1397,7 @@ pub async fn reload_worker_config( } } -pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { +pub async fn load_base_url(db: &DB) -> error::Result { let q_base_url = load_value_from_global_settings(db, BASE_URL_SETTING).await?; let std_base_url = std::env::var("BASE_URL") @@ -1421,6 +1421,14 @@ pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { std_base_url }; + { + let mut l = BASE_URL.write().await; + *l = base_url.clone(); + } + Ok(base_url) +} + +pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { let q_oauth = load_value_from_global_settings(db, OAUTH_SETTING).await?; let oauths = if let Some(q) = q_oauth { @@ -1434,6 +1442,7 @@ pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { None }; + let base_url = load_base_url(db).await?; let is_secure = base_url.starts_with("https://"); { @@ -1443,11 +1452,6 @@ pub async fn reload_base_url_setting(db: &DB) -> error::Result<()> { .unwrap(); } - { - let mut l = BASE_URL.write().await; - *l = base_url - } - { let mut l = IS_SECURE.write().await; *l = is_secure; diff --git a/backend/tests/worker.rs b/backend/tests/worker.rs index 2aebf6f5e8..571b71cd8f 100644 --- a/backend/tests/worker.rs +++ b/backend/tests/worker.rs @@ -80,6 +80,7 @@ async fn initialize_tracing() { let _ = windmill_common::tracing_init::initialize_tracing( "test", &windmill_common::utils::Mode::Standalone, + "test", ); }); } diff --git a/backend/windmill-common/src/otel_ee.rs b/backend/windmill-common/src/otel_ee.rs index e33f2ff486..f3ada162f6 100644 --- a/backend/windmill-common/src/otel_ee.rs +++ b/backend/windmill-common/src/otel_ee.rs @@ -38,7 +38,7 @@ pub trait FutureExt: Sized { use tracing_subscriber::EnvFilter; -pub(crate) fn init_logs_bridge(_mode: &Mode, _hostname: &str) -> Option { +pub(crate) fn init_logs_bridge(_mode: &Mode, _hostname: &str, _env: &str) -> Option { None } @@ -46,11 +46,12 @@ pub(crate) fn init_logs_bridge(_mode: &Mode, _hostname: &str) -> Option Option { None } -pub(crate) fn init_meter_provider(_mode: &Mode, _hostname: &str) -> OtelProvider { +pub(crate) fn init_meter_provider(_mode: &Mode, _hostname: &str, _env: &str) -> OtelProvider { None } diff --git a/backend/windmill-common/src/tracing_init.rs b/backend/windmill-common/src/tracing_init.rs index 962f99204e..aa1e409dd1 100644 --- a/backend/windmill-common/src/tracing_init.rs +++ b/backend/windmill-common/src/tracing_init.rs @@ -47,6 +47,7 @@ pub const TMP_WINDMILL_LOGS_SERVICE: &str = concatcp!("/tmp/windmill/", LOGS_SER pub fn initialize_tracing( hostname: &str, mode: &Mode, + environment: &str, ) -> (WorkerGuard, crate::otel_ee::OtelProvider) { let style = std::env::var("RUST_LOG_STYLE").unwrap_or_else(|_| "auto".into()); @@ -57,16 +58,16 @@ pub fn initialize_tracing( ) } - let meter_provider = crate::otel_ee::init_meter_provider(mode, hostname); + let meter_provider = crate::otel_ee::init_meter_provider(mode, hostname, environment); #[cfg(all(feature = "otel", feature = "enterprise"))] - let opentelemetry = crate::otel_ee::init_otlp_tracer(mode, hostname) + let opentelemetry = crate::otel_ee::init_otlp_tracer(mode, hostname, environment) .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, hostname); + let logs_bridge = crate::otel_ee::init_logs_bridge(&mode, hostname, environment); use tracing_appender::rolling::{RollingFileAppender, Rotation}; diff --git a/backend/windmill-worker/src/handle_child.rs b/backend/windmill-worker/src/handle_child.rs index cb35070a40..54eda3b1ba 100644 --- a/backend/windmill-worker/src/handle_child.rs +++ b/backend/windmill-worker/src/handle_child.rs @@ -126,7 +126,7 @@ pub async fn handle_child( let (tx, rx) = broadcast::channel::<()>(3); let mut rx2 = tx.subscribe(); - let output = child_joined_output_stream(&mut child); + let output = child_joined_output_stream(&mut child, job_id.clone()); let job_id = job_id.clone(); @@ -677,6 +677,7 @@ where /// builds a stream joining both stdout and stderr each read line by line fn child_joined_output_stream( child: &mut Child, + job_id: Uuid, ) -> impl stream::FusedStream> { let stderr = child .stderr @@ -691,19 +692,20 @@ 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), + lines_to_stream(stderr, true, job_id.clone()), + lines_to_stream(stdout, false, job_id), ) } pub fn lines_to_stream( mut lines: tokio::io::Lines, stderr: bool, + job_id: Uuid, ) -> impl futures::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, &job_id)) }) } diff --git a/backend/windmill-worker/src/job_logger_ee.rs b/backend/windmill-worker/src/job_logger_ee.rs index 5414ccd496..310419e3d1 100644 --- a/backend/windmill-worker/src/job_logger_ee.rs +++ b/backend/windmill-worker/src/job_logger_ee.rs @@ -34,6 +34,7 @@ pub(crate) async fn default_disk_log_storage( pub(crate) fn process_streaming_log_lines( r: Result, io::Error>, _stderr: bool, + _job_id: &Uuid, ) -> Option> { r.transpose() }