diff --git a/backend/windmill-worker/src/bash_executor.rs b/backend/windmill-worker/src/bash_executor.rs index 429fd88c94..7d92607af5 100644 --- a/backend/windmill-worker/src/bash_executor.rs +++ b/backend/windmill-worker/src/bash_executor.rs @@ -215,7 +215,10 @@ exit $exit_status .current_dir(job_dir) .env_clear() .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Bash).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Bash, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) .args(cmd_args) @@ -241,7 +244,10 @@ exit $exit_status .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Bash).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Bash, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) .env("HOME", HOME_ENV.as_str()) diff --git a/backend/windmill-worker/src/bun_executor.rs b/backend/windmill-worker/src/bun_executor.rs index 914c3ca57b..6ada66c4bb 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -1564,7 +1564,9 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Bun, &job.id, &job.workspace_id, conn).await?, + ) .envs(common_bun_proc_envs) .env("PATH", PATH_ENV.as_str()) .args(args) @@ -1582,7 +1584,10 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Bun, &job.id, &job.workspace_id, conn) + .await?, + ) .envs(common_bun_proc_envs) .stdin(Stdio::null()) .stdout(Stdio::piped()) @@ -1613,7 +1618,10 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Bun, &job.id, &job.workspace_id, conn) + .await?, + ) .envs(common_bun_proc_envs) .stdin(Stdio::null()) .stdout(Stdio::piped()) diff --git a/backend/windmill-worker/src/csharp_executor.rs b/backend/windmill-worker/src/csharp_executor.rs index 4c3e63f4c9..22f46e724b 100644 --- a/backend/windmill-worker/src/csharp_executor.rs +++ b/backend/windmill-worker/src/csharp_executor.rs @@ -600,7 +600,10 @@ pub async fn handle_csharp_job( .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::CSharp).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::CSharp, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -633,7 +636,10 @@ pub async fn handle_csharp_job( .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::CSharp).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::CSharp, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("DOTNET_CLI_HOME", &*CSHARP_CACHE_DIR) diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index 6cbf85c47b..d9f6785eb6 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -121,11 +121,13 @@ async fn get_common_deno_proc_envs( } // Add proxy envs (including OTEL tracing proxy if enabled for deno) - for (k, v) in get_proxy_envs_for_lang(&ScriptLang::Deno) - .await - .unwrap_or_default() - { - deno_envs.insert(k.to_string(), v); + if let Some(conn) = conn { + for (k, v) in get_proxy_envs_for_lang(&ScriptLang::Deno, job_id, w_id, conn) + .await + .unwrap_or_default() + { + deno_envs.insert(k.to_string(), v); + } } return deno_envs; diff --git a/backend/windmill-worker/src/go_executor.rs b/backend/windmill-worker/src/go_executor.rs index ddcc108659..0315bf25c3 100644 --- a/backend/windmill-worker/src/go_executor.rs +++ b/backend/windmill-worker/src/go_executor.rs @@ -354,7 +354,7 @@ func Run(req Req) (interface{{}}, error){{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Go).await?) + .envs(get_proxy_envs_for_lang(&ScriptLang::Go, &job.id, &job.workspace_id, conn).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -375,7 +375,7 @@ func Run(req Req) (interface{{}}, error){{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Go).await?) + .envs(get_proxy_envs_for_lang(&ScriptLang::Go, &job.id, &job.workspace_id, conn).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) diff --git a/backend/windmill-worker/src/nu_executor.rs b/backend/windmill-worker/src/nu_executor.rs index ca309105b5..28ac27c925 100644 --- a/backend/windmill-worker/src/nu_executor.rs +++ b/backend/windmill-worker/src/nu_executor.rs @@ -264,7 +264,7 @@ async fn run<'a>( .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Nu).await?) + .envs(get_proxy_envs_for_lang(&ScriptLang::Nu, &job.id, &job.workspace_id, conn).await?) .args(vec![ "--config", "run.config.proto", @@ -303,7 +303,7 @@ async fn run<'a>( .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Nu).await?) + .envs(get_proxy_envs_for_lang(&ScriptLang::Nu, &job.id, &job.workspace_id, conn).await?) // TODO(v1): // "--plugins", // &format!( diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index 619087659d..320bcb3e40 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -841,7 +841,10 @@ mount {{ .env_clear() // inject PYTHONPATH here - for some reason I had to do it in nsjail conf .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Python3).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Python3, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -867,7 +870,10 @@ mount {{ .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Python3).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Python3, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) diff --git a/backend/windmill-worker/src/ruby_executor.rs b/backend/windmill-worker/src/ruby_executor.rs index 11031fa481..d283b2a1c4 100644 --- a/backend/windmill-worker/src/ruby_executor.rs +++ b/backend/windmill-worker/src/ruby_executor.rs @@ -812,7 +812,10 @@ mount {{ .envs(envs) .envs(reserved_variables) .envs(RUBY_PROXY_ENVS.clone()) - .envs(get_proxy_envs_for_lang(&ScriptLang::Ruby).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Ruby, &job.id, &job.workspace_id, conn) + .await?, + ) .args(vec![ "--config", "run.config.proto", @@ -851,7 +854,10 @@ mount {{ .env("BASE_INTERNAL_URL", base_internal_url) .envs(reserved_variables) .envs(RUBY_PROXY_ENVS.clone()) - .envs(get_proxy_envs_for_lang(&ScriptLang::Ruby).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Ruby, &job.id, &job.workspace_id, conn) + .await?, + ) .envs(envs); cmd.stdin(Stdio::null()) diff --git a/backend/windmill-worker/src/rust_executor.rs b/backend/windmill-worker/src/rust_executor.rs index a1f29696c4..dae8e08766 100644 --- a/backend/windmill-worker/src/rust_executor.rs +++ b/backend/windmill-worker/src/rust_executor.rs @@ -700,7 +700,10 @@ pub async fn handle_rust_job( .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Rust).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Rust, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -716,7 +719,10 @@ pub async fn handle_rust_job( .env_clear() .envs(envs) .envs(reserved_variables) - .envs(get_proxy_envs_for_lang(&ScriptLang::Rust).await?) + .envs( + get_proxy_envs_for_lang(&ScriptLang::Rust, &job.id, &job.workspace_id, conn) + .await?, + ) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 4202064201..934b10f8f4 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -745,21 +745,37 @@ pub async fn is_otel_tracing_proxy_enabled_for_lang(lang: &ScriptLang) -> bool { /// Otherwise, uses the standard HTTP_PROXY/HTTPS_PROXY from environment. pub async fn get_proxy_envs_for_lang( lang: &ScriptLang, + job_id: &uuid::Uuid, + w_id: &str, + conn: &Connection, ) -> anyhow::Result> { #[cfg(all(feature = "private", feature = "enterprise"))] if is_otel_tracing_proxy_enabled_for_lang(lang).await { - return get_otel_tracing_proxy_envs().await; + return get_otel_tracing_proxy_envs(job_id, w_id, conn).await; } - let _ = lang; + let _ = (lang, job_id, w_id, conn); Ok(PROXY_ENVS.clone()) } #[cfg(all(feature = "private", feature = "enterprise"))] -async fn get_otel_tracing_proxy_envs() -> anyhow::Result> { - let port = crate::otel_tracing_proxy_ee::TRACING_PROXY_PORT +async fn get_otel_tracing_proxy_envs( + job_id: &uuid::Uuid, + w_id: &str, + conn: &Connection, +) -> anyhow::Result> { + let port = match *crate::otel_tracing_proxy_ee::TRACING_PROXY_PORT .read() .await - .ok_or_else(|| anyhow::anyhow!("OTEL tracing proxy port not initialized"))?; + { + Some(p) => p, + None => { + let reason = "OTEL tracing proxy is enabled but not available (not initialized yet, or NUM_WORKERS > 1). \ + This job's HTTP requests will not be traced."; + tracing::warn!("{}", reason); + append_logs(job_id, w_id, format!("\n[warning] {reason}\n"), conn).await; + return Ok(PROXY_ENVS.clone()); + } + }; let proxy_url = format!("http://127.0.0.1:{}", port); Ok(vec![ ("HTTP_PROXY", proxy_url.clone()),