From 95df7b9a6a8ffcbca92b3249a61d97c32c9dbc4f Mon Sep 17 00:00:00 2001 From: Pyra <92104930+pyranota@users.noreply.github.com> Date: Tue, 20 Jan 2026 14:38:02 +0100 Subject: [PATCH] feat: otel REST tracing (#7571) --- ...7a356217bbae06b8453f0c32046f0cbf20dcb.json | 23 + ...028803a286e69443a151cd1e3d16da7e8e4f7.json | 29 ++ ...1e4f15a825b442346c7386836127bc5723c93.json | 22 + backend/Cargo.lock | 396 ++++++++++++++++-- backend/Cargo.toml | 8 + backend/ee-repo-ref.txt | 2 +- .../20260112000000_otel_traces.down.sql | 3 + .../20260112000000_otel_traces.up.sql | 36 ++ backend/src/main.rs | 57 ++- backend/src/monitor.rs | 40 +- backend/windmill-api/openapi.yaml | 20 + backend/windmill-api/src/jobs.rs | 61 +++ backend/windmill-common/Cargo.toml | 2 +- .../windmill-common/src/global_settings.rs | 1 + backend/windmill-worker/Cargo.toml | 10 +- .../nsjail/download.py.config.proto | 9 + .../nsjail/download.ruby.config.proto | 9 +- .../nsjail/download.rust.config.proto | 9 +- .../nsjail/lock.ruby.config.proto | 9 +- .../nsjail/run.bash.config.proto | 9 + .../nsjail/run.bun.config.proto | 7 + .../nsjail/run.csharp.config.proto | 9 + .../nsjail/run.go.config.proto | 7 + .../nsjail/run.nu.config.proto | 9 + .../nsjail/run.python3.config.proto | 7 + .../nsjail/run.ruby.config.proto | 8 +- .../nsjail/run.rust.config.proto | 9 +- backend/windmill-worker/src/bash_executor.rs | 18 +- backend/windmill-worker/src/bun_executor.rs | 13 +- backend/windmill-worker/src/common.rs | 14 + .../windmill-worker/src/csharp_executor.rs | 14 +- backend/windmill-worker/src/deno_executor.rs | 10 +- backend/windmill-worker/src/go_executor.rs | 13 +- backend/windmill-worker/src/handle_child.rs | 5 + backend/windmill-worker/src/java_executor.rs | 1 - backend/windmill-worker/src/lib.rs | 6 + backend/windmill-worker/src/nu_executor.rs | 15 +- .../src/otel_tracing_proxy_oss.rs | 13 + .../windmill-worker/src/python_executor.rs | 14 +- backend/windmill-worker/src/ruby_executor.rs | 32 +- backend/windmill-worker/src/rust_executor.rs | 28 +- backend/windmill-worker/src/worker.rs | 67 ++- flake.nix | 1 + .../components/FlowStatusViewerInner.svelte | 13 +- .../src/lib/components/InstanceSetting.svelte | 35 ++ .../src/lib/components/JobOtelTraces.svelte | 311 ++++++++++++++ .../src/lib/components/instanceSettings.ts | 25 +- .../components/scriptEditor/LogPanel.svelte | 11 + .../(root)/(logged)/run/[...run]/+page.svelte | 8 +- 49 files changed, 1357 insertions(+), 121 deletions(-) create mode 100644 backend/.sqlx/query-002d68d7c4437522a6dae95af007a356217bbae06b8453f0c32046f0cbf20dcb.json create mode 100644 backend/.sqlx/query-5b5cb5339208847bd542f8d903d028803a286e69443a151cd1e3d16da7e8e4f7.json create mode 100644 backend/.sqlx/query-90d93fd3bd91e468c1e796e41e31e4f15a825b442346c7386836127bc5723c93.json create mode 100644 backend/migrations/20260112000000_otel_traces.down.sql create mode 100644 backend/migrations/20260112000000_otel_traces.up.sql create mode 100644 backend/windmill-worker/src/otel_tracing_proxy_oss.rs create mode 100644 frontend/src/lib/components/JobOtelTraces.svelte diff --git a/backend/.sqlx/query-002d68d7c4437522a6dae95af007a356217bbae06b8453f0c32046f0cbf20dcb.json b/backend/.sqlx/query-002d68d7c4437522a6dae95af007a356217bbae06b8453f0c32046f0cbf20dcb.json new file mode 100644 index 0000000000..56616334ac --- /dev/null +++ b/backend/.sqlx/query-002d68d7c4437522a6dae95af007a356217bbae06b8453f0c32046f0cbf20dcb.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT created_by FROM v2_job WHERE id = $1 AND workspace_id = $2", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "created_by", + "type_info": "Varchar" + } + ], + "parameters": { + "Left": [ + "Uuid", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "002d68d7c4437522a6dae95af007a356217bbae06b8453f0c32046f0cbf20dcb" +} diff --git a/backend/.sqlx/query-5b5cb5339208847bd542f8d903d028803a286e69443a151cd1e3d16da7e8e4f7.json b/backend/.sqlx/query-5b5cb5339208847bd542f8d903d028803a286e69443a151cd1e3d16da7e8e4f7.json new file mode 100644 index 0000000000..4a1aceb5f1 --- /dev/null +++ b/backend/.sqlx/query-5b5cb5339208847bd542f8d903d028803a286e69443a151cd1e3d16da7e8e4f7.json @@ -0,0 +1,29 @@ +{ + "db_name": "PostgreSQL", + "query": "INSERT INTO otel_traces (\n trace_id, span_id, trace_state, parent_span_id, flags,\n name, kind, start_time_unix_nano, end_time_unix_nano,\n attributes, dropped_attributes_count,\n events, dropped_events_count,\n links, dropped_links_count,\n status\n ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Bytea", + "Bytea", + "Text", + "Bytea", + "Int4", + "Text", + "Int4", + "Int8", + "Int8", + "Jsonb", + "Int4", + "Jsonb", + "Int4", + "Jsonb", + "Int4", + "Jsonb" + ] + }, + "nullable": [] + }, + "hash": "5b5cb5339208847bd542f8d903d028803a286e69443a151cd1e3d16da7e8e4f7" +} diff --git a/backend/.sqlx/query-90d93fd3bd91e468c1e796e41e31e4f15a825b442346c7386836127bc5723c93.json b/backend/.sqlx/query-90d93fd3bd91e468c1e796e41e31e4f15a825b442346c7386836127bc5723c93.json new file mode 100644 index 0000000000..5f27525418 --- /dev/null +++ b/backend/.sqlx/query-90d93fd3bd91e468c1e796e41e31e4f15a825b442346c7386836127bc5723c93.json @@ -0,0 +1,22 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT json_build_object(\n 'trace_id', encode(trace_id, 'hex'), -- BYTEA to hex string\n 'span_id', encode(span_id, 'hex'), -- BYTEA to hex string\n 'parent_span_id', encode(parent_span_id, 'hex'), -- BYTEA to hex string\n 'trace_state', trace_state,\n 'flags', flags,\n 'name', name,\n 'kind', kind,\n 'start_time_unix_nano', start_time_unix_nano,\n 'end_time_unix_nano', end_time_unix_nano,\n 'attributes', attributes,\n 'dropped_attributes_count', dropped_attributes_count,\n 'events', events,\n 'dropped_events_count', dropped_events_count,\n 'links', links,\n 'dropped_links_count', dropped_links_count,\n 'status', status\n ) as \"span!\"\n FROM otel_traces\n WHERE trace_id = $1\n ORDER BY start_time_unix_nano", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "span!", + "type_info": "Json" + } + ], + "parameters": { + "Left": [ + "Bytea" + ] + }, + "nullable": [ + null + ] + }, + "hash": "90d93fd3bd91e468c1e796e41e31e4f15a825b442346c7386836127bc5723c93" +} diff --git a/backend/Cargo.lock b/backend/Cargo.lock index 8386ade173..6629d852e6 100644 --- a/backend/Cargo.lock +++ b/backend/Cargo.lock @@ -517,8 +517,24 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7f6fd5ddaf0351dff5b8da21b2fb4ff8e08ddd02857f0bf69c47639106c0fff0" dependencies = [ - "asn1-rs-derive", - "asn1-rs-impl", + "asn1-rs-derive 0.4.0", + "asn1-rs-impl 0.1.0", + "displaydoc", + "nom 7.1.3", + "num-traits", + "rusticata-macros", + "thiserror 1.0.69", + "time", +] + +[[package]] +name = "asn1-rs" +version = "0.6.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5493c3bedbacf7fd7382c6346bbd66687d12bbaad3a89a2d2c303ee6cf20b048" +dependencies = [ + "asn1-rs-derive 0.5.1", + "asn1-rs-impl 0.2.0", "displaydoc", "nom 7.1.3", "num-traits", @@ -539,6 +555,18 @@ dependencies = [ "synstructure 0.12.6", ] +[[package]] +name = "asn1-rs-derive" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "965c2d33e53cb6b267e148a4cb0760bc01f4904c1cd4bb4002a085bb016d1490" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.114", + "synstructure 0.13.2", +] + [[package]] name = "asn1-rs-impl" version = "0.1.0" @@ -550,6 +578,17 @@ dependencies = [ "syn 1.0.109", ] +[[package]] +name = "asn1-rs-impl" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7b18050c2cd6fe86c3a76584ef5e0baf286d038cda203eb6223df2cc413565f7" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.114", +] + [[package]] name = "ast_node" version = "0.9.9" @@ -620,6 +659,17 @@ dependencies = [ "zstd-safe", ] +[[package]] +name = "async-lock" +version = "3.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "290f7f2596bd5b78a9fec8088ccd89180d7f9f55b94b0576823bbbdc72ee8311" +dependencies = [ + "event-listener 5.4.1", + "event-listener-strategy", + "pin-project-lite", +] + [[package]] name = "async-nats" version = "0.38.0" @@ -1853,6 +1903,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63044e1ae8e69f3b5a92c736ca6269b8d12fa7efe39bf34ddb06d102cf0e2cab" dependencies = [ "memchr", + "regex-automata", "serde", ] @@ -4148,7 +4199,7 @@ dependencies = [ "winapi", "windows-sys 0.59.0", "x25519-dalek", - "x509-parser", + "x509-parser 0.15.1", "yoke 0.7.5", ] @@ -4429,11 +4480,11 @@ dependencies = [ "hyper-util", "log", "once_cell", - "opentelemetry", + "opentelemetry 0.27.1", "opentelemetry-http", "opentelemetry-otlp", "opentelemetry-semantic-conventions", - "opentelemetry_sdk", + "opentelemetry_sdk 0.27.1", "pin-project", "serde", "thiserror 2.0.18", @@ -4671,7 +4722,21 @@ version = "8.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "dbd676fbbab537128ef0278adb5576cf363cff6aa22a7b24effe97347cfab61e" dependencies = [ - "asn1-rs", + "asn1-rs 0.5.2", + "displaydoc", + "nom 7.1.3", + "num-bigint", + "num-traits", + "rusticata-macros", +] + +[[package]] +name = "der-parser" +version = "9.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5cd0a5c643689626bec213c4d8bd4d96acc8ffdb4ad4bb6bc16abf27d5f4b553" +dependencies = [ + "asn1-rs 0.6.2", "displaydoc", "nom 7.1.3", "num-bigint", @@ -4755,13 +4820,33 @@ dependencies = [ "syn 2.0.114", ] +[[package]] +name = "derive_more" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4a9b99b9cbbe49445b21764dc0625032a89b145a2642e67603e1c936f5458d05" +dependencies = [ + "derive_more-impl 1.0.0", +] + [[package]] name = "derive_more" version = "2.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d751e9e49156b02b44f9c1815bcb94b984cdcc4396ecc32521c739452808b134" dependencies = [ - "derive_more-impl", + "derive_more-impl 2.1.1", +] + +[[package]] +name = "derive_more-impl" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb7330aeadfbe296029522e6c40f315320aba36fc43a5b3632f3795348f3bd22" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.114", ] [[package]] @@ -6045,6 +6130,21 @@ dependencies = [ "seq-macro", ] +[[package]] +name = "generator" +version = "0.8.8" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "52f04ae4152da20c76fe800fa48659201d5cf627c5149ca0b707b69d7eef6cf9" +dependencies = [ + "cc", + "cfg-if", + "libc", + "log", + "rustversion", + "windows-link 0.2.1", + "windows-result 0.4.1", +] + [[package]] name = "generic-array" version = "0.14.7" @@ -6764,6 +6864,35 @@ version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" +[[package]] +name = "hudsucker" +version = "0.22.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8bb9d62508d54891fe529dc3a3e169aa7938b89898ba5ab0431ac5bafe66a249" +dependencies = [ + "async-compression", + "bstr", + "futures", + "http 1.4.0", + "http-body-util", + "hyper 1.8.1", + "hyper-rustls 0.26.0", + "hyper-tls", + "hyper-tungstenite", + "hyper-util", + "moka", + "rand 0.8.5", + "rcgen", + "thiserror 1.0.69", + "time", + "tokio", + "tokio-graceful", + "tokio-rustls 0.25.0", + "tokio-tungstenite 0.21.0", + "tokio-util", + "tracing", +] + [[package]] name = "humantime" version = "2.3.0" @@ -6829,10 +6958,13 @@ dependencies = [ "http 1.4.0", "hyper 1.8.1", "hyper-rustls 0.27.7", + "hyper-tls", "hyper-util", + "native-tls", "pin-project-lite", "rustls-native-certs 0.7.3", "tokio", + "tokio-native-tls", "tokio-rustls 0.26.4", "tower-service", ] @@ -6868,6 +7000,25 @@ dependencies = [ "tokio-rustls 0.24.1", ] +[[package]] +name = "hyper-rustls" +version = "0.26.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a0bea761b46ae2b24eb4aef630d8d1c398157b6fc29e6350ecf090a0b70c952c" +dependencies = [ + "futures-util", + "http 1.4.0", + "hyper 1.8.1", + "hyper-util", + "log", + "rustls 0.22.4", + "rustls-pki-types", + "tokio", + "tokio-rustls 0.25.0", + "tower-service", + "webpki-roots 0.26.11", +] + [[package]] name = "hyper-rustls" version = "0.27.7" @@ -6916,6 +7067,21 @@ dependencies = [ "tower-service", ] +[[package]] +name = "hyper-tungstenite" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7a343d17fe7885302ed7252767dc7bb83609a874b6ff581142241ec4b73957ad" +dependencies = [ + "http-body-util", + "hyper 1.8.1", + "hyper-util", + "pin-project-lite", + "tokio", + "tokio-tungstenite 0.21.0", + "tungstenite 0.21.0", +] + [[package]] name = "hyper-util" version = "0.1.19" @@ -7953,6 +8119,22 @@ version = "0.4.29" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897" +[[package]] +name = "loom" +version = "0.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "419e0dc8046cb947daa77eb95ae174acfbddb7673b4151f56d1eed8e93fbfaca" +dependencies = [ + "cfg-if", + "generator", + "pin-utils", + "scoped-tls", + "serde", + "serde_json", + "tracing", + "tracing-subscriber", +] + [[package]] name = "lru" version = "0.12.5" @@ -8387,10 +8569,13 @@ version = "0.12.12" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a3dec6bd31b08944e08b58fd99373893a6c17054d6f3ea5006cc894f4f4eee2a" dependencies = [ + "async-lock", "crossbeam-channel", "crossbeam-epoch", "crossbeam-utils", "equivalent", + "event-listener 5.4.1", + "futures-util", "parking_lot", "portable-atomic", "smallvec", @@ -9143,7 +9328,16 @@ version = "0.6.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9bedf36ffb6ba96c2eb7144ef6270557b52e54b20c0a8e1eb2ff99a6c6959bff" dependencies = [ - "asn1-rs", + "asn1-rs 0.5.2", +] + +[[package]] +name = "oid-registry" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a8d8034d9489cdaf79228eb9f6a3b8d7bb32ba00d6645ebd48eef4077ceb5bd9" +dependencies = [ + "asn1-rs 0.6.2", ] [[package]] @@ -9303,13 +9497,27 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry" +version = "0.29.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9e87237e2775f74896f9ad219d26a2081751187eb7c9f5c58dde20a23b95d16c" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror 2.0.17", + "tracing", +] + [[package]] name = "opentelemetry-appender-tracing" version = "0.27.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ab5feffc321035ad94088a7e5333abb4d84a8726e54a802e736ce9dd7237e85b" dependencies = [ - "opentelemetry", + "opentelemetry 0.27.1", "tracing", "tracing-core", "tracing-subscriber", @@ -9324,7 +9532,7 @@ dependencies = [ "async-trait", "bytes", "http 1.4.0", - "opentelemetry", + "opentelemetry 0.27.1", ] [[package]] @@ -9336,10 +9544,10 @@ dependencies = [ "async-trait", "futures-core", "http 1.4.0", - "opentelemetry", + "opentelemetry 0.27.1", "opentelemetry-http", - "opentelemetry-proto", - "opentelemetry_sdk", + "opentelemetry-proto 0.27.0", + "opentelemetry_sdk 0.27.1", "prost", "serde_json", "thiserror 1.0.69", @@ -9355,13 +9563,29 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6e05acbfada5ec79023c85368af14abd0b307c015e9064d249b2a950ef459a6" dependencies = [ "hex", - "opentelemetry", - "opentelemetry_sdk", + "opentelemetry 0.27.1", + "opentelemetry_sdk 0.27.1", "prost", "serde", "tonic", ] +[[package]] +name = "opentelemetry-proto" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8c40da242381435e18570d5b9d50aca2a4f4f4d8e146231adb4e7768023309b3" +dependencies = [ + "base64 0.22.1", + "hex", + "opentelemetry 0.29.1", + "opentelemetry_sdk 0.29.0", + "prost", + "serde", + "tonic", + "tracing", +] + [[package]] name = "opentelemetry-semantic-conventions" version = "0.27.0" @@ -9379,7 +9603,7 @@ dependencies = [ "futures-executor", "futures-util", "glob", - "opentelemetry", + "opentelemetry 0.27.1", "percent-encoding", "rand 0.8.5", "serde_json", @@ -9389,6 +9613,23 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry_sdk" +version = "0.29.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "afdefb21d1d47394abc1ba6c57363ab141be19e27cc70d0e422b7f303e4d290b" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "glob", + "opentelemetry 0.29.1", + "percent-encoding", + "rand 0.9.0", + "serde_json", + "thiserror 2.0.17", +] + [[package]] name = "option-ext" version = "0.2.0" @@ -10600,6 +10841,20 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "rcgen" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "75e669e5202259b5314d1ea5397316ad400819437857b90861765f24c4cf80a2" +dependencies = [ + "pem 3.0.5", + "ring 0.17.14", + "rustls-pki-types", + "time", + "x509-parser 0.16.0", + "yasna", +] + [[package]] name = "rdkafka" version = "0.36.2" @@ -13780,6 +14035,19 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-graceful" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "627ba4daa4cbce14740603401c895e72d47ecd86690a18e3f0841266e9340de7" +dependencies = [ + "loom", + "pin-project-lite", + "slab", + "tokio", + "tracing", +] + [[package]] name = "tokio-macros" version = "2.5.0" @@ -13928,6 +14196,24 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-tungstenite" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c83b561d025642014097b66e6c1bb422783339e0909e4429cde4749d1990bc38" +dependencies = [ + "futures-util", + "log", + "native-tls", + "rustls 0.22.4", + "rustls-pki-types", + "tokio", + "tokio-native-tls", + "tokio-rustls 0.25.0", + "tungstenite 0.21.0", + "webpki-roots 0.26.11", +] + [[package]] name = "tokio-tungstenite" version = "0.24.0" @@ -13939,7 +14225,7 @@ dependencies = [ "native-tls", "tokio", "tokio-native-tls", - "tungstenite", + "tungstenite 0.24.0", ] [[package]] @@ -14245,8 +14531,8 @@ checksum = "97a971f6058498b5c0f1affa23e7ea202057a7301dbff68e968b2d578bcbd053" dependencies = [ "js-sys", "once_cell", - "opentelemetry", - "opentelemetry_sdk", + "opentelemetry 0.27.1", + "opentelemetry_sdk 0.27.1", "smallvec", "tracing", "tracing-core", @@ -14360,6 +14646,28 @@ dependencies = [ "tokio", ] +[[package]] +name = "tungstenite" +version = "0.21.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9ef1a641ea34f399a848dea702823bbecfb4c486f911735368f1f137cb8257e1" +dependencies = [ + "byteorder", + "bytes", + "data-encoding", + "http 1.4.0", + "httparse", + "log", + "native-tls", + "rand 0.8.5", + "rustls 0.22.4", + "rustls-pki-types", + "sha1", + "thiserror 1.0.69", + "url", + "utf-8", +] + [[package]] name = "tungstenite" version = "0.24.0" @@ -15265,6 +15573,7 @@ dependencies = [ "memchr", "object_store", "once_cell", + "opentelemetry-proto 0.29.0", "pep440_rs", "prometheus", "quote", @@ -15401,7 +15710,7 @@ dependencies = [ "tokio-postgres 0.7.11", "tokio-postgres 0.7.13", "tokio-stream", - "tokio-tungstenite", + "tokio-tungstenite 0.24.0", "tokio-util", "tonic", "tower 0.5.3", @@ -15515,11 +15824,11 @@ dependencies = [ "object_store", "once_cell", "openidconnect", - "opentelemetry", + "opentelemetry 0.27.1", "opentelemetry-appender-tracing", "opentelemetry-otlp", "opentelemetry-semantic-conventions", - "opentelemetry_sdk", + "opentelemetry_sdk 0.27.1", "pep440_rs", "phf 0.11.3", "pin-project-lite", @@ -15949,6 +16258,7 @@ dependencies = [ "deno_url", "deno_web", "deno_webidl", + "derive_more 1.0.0", "dotenv", "dyn-iter", "eventsource-stream", @@ -15957,6 +16267,10 @@ dependencies = [ "gcp_auth", "git-version", "hex", + "hudsucker", + "hyper-http-proxy", + "hyper-tls", + "hyper-util", "itertools 0.14.0", "jsonwebtoken 8.3.0", "lazy_static", @@ -15969,7 +16283,8 @@ dependencies = [ "nix 0.27.1", "object_store", "once_cell", - "opentelemetry", + "opentelemetry 0.27.1", + "opentelemetry-proto 0.29.0", "oracle", "pem 3.0.5", "pep440_rs", @@ -15977,6 +16292,7 @@ dependencies = [ "process-wrap", "prometheus", "rand 0.9.0", + "rcgen", "regex", "reqwest 0.13.1", "reqwest-middleware", @@ -15992,6 +16308,7 @@ dependencies = [ "tokio-stream", "tokio-util", "tracing", + "tracing-opentelemetry", "ulid", "url", "urlencoding", @@ -16673,12 +16990,30 @@ version = "0.15.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7069fba5b66b9193bd2c5d3d4ff12b839118f6bcbef5328efafafb5395cf63da" dependencies = [ - "asn1-rs", + "asn1-rs 0.5.2", "data-encoding", - "der-parser", + "der-parser 8.2.0", "lazy_static", "nom 7.1.3", - "oid-registry", + "oid-registry 0.6.1", + "rusticata-macros", + "thiserror 1.0.69", + "time", +] + +[[package]] +name = "x509-parser" +version = "0.16.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fcbc162f30700d6f3f82a24bf7cc62ffe7caea42c0b2cba8bf7f3ae50cf51f69" +dependencies = [ + "asn1-rs 0.6.2", + "data-encoding", + "der-parser 9.0.0", + "lazy_static", + "nom 7.1.3", + "oid-registry 0.7.1", + "ring 0.17.14", "rusticata-macros", "thiserror 1.0.69", "time", @@ -16730,6 +17065,15 @@ version = "0.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "09041cd90cf85f7f8b2df60c646f853b7f535ce68f85244eb6731cf89fa498ec" +[[package]] +name = "yasna" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e17bb3549cc1321ae1296b9cdc2698e2b6cb1992adfa19a8c72e5b7a738f44cd" +dependencies = [ + "time", +] + [[package]] name = "yoke" version = "0.7.5" diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 9828a127a7..1a785a2ebd 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -161,6 +161,7 @@ k8s-openapi.workspace = true libloading.workspace = true bitflags.workspace = true globset.workspace = true +opentelemetry-proto.workspace = true [target.'cfg(windows)'.dependencies] @@ -219,6 +220,8 @@ memchr = "2.7.4" axum = { version = "^0.7", features = ["multipart", "macros"] } headers = "^0" hyper = { version = "^1", features = ["full"] } +hyper-tls = "^0.6" +hyper-util = { version = "^0.1", features = ["client-legacy", "http1", "tokio"] } tokio = { version = "=1.46.1", features = ["full", "tracing", "time"] } tokio-stream = { version = "0.1.17" } tower = "^0" @@ -233,6 +236,7 @@ thiserror = "^2" anyhow = "^1" chrono = { version = "^0.4", features = ["serde"] } chrono-tz = "^0.10.1" +derive_more = { version = "1", features = ["deref", "deref_mut"], default-features = false } tracing = "^0" tracing-subscriber = { version = "^0", features = ["env-filter", "json"] } tracing-appender = "^0" @@ -415,6 +419,7 @@ opentelemetry_sdk = { version = "0.27.1", features = ["rt-tokio"] } opentelemetry-otlp = { version = "0.27.0", features = ["grpc-tonic", "tls"] } opentelemetry-appender-tracing = "0.27.0" opentelemetry-semantic-conventions = { version = "0.27.0", features = ["semconv_experimental"] } +opentelemetry-proto = { version = "0.29.0", features = ["with-serde", "gen-tonic"] } bollard = "0.18.1" @@ -456,3 +461,6 @@ oracle = { version = "0.6.3", features = ["chrono"] } rumqttc = { version = "0.24.0", features = ["use-native-tls"]} strum = { version = "0.27", features = ["derive"] } strum_macros = "^0" +hudsucker = { version = "0.22", features = ["rcgen-ca", "native-tls-client"] } +hyper-http-proxy = { version = "1", default-features = false, features = ["native-tls"] } +rcgen = "0.13" diff --git a/backend/ee-repo-ref.txt b/backend/ee-repo-ref.txt index 6bfbf852e9..f744b6679b 100644 --- a/backend/ee-repo-ref.txt +++ b/backend/ee-repo-ref.txt @@ -1 +1 @@ -62cc5aaf46a4f825d9b0cf446924c07eeb95b8d5 +0056bec3ccc0d7a80461a2d036c555364d30b0d6 diff --git a/backend/migrations/20260112000000_otel_traces.down.sql b/backend/migrations/20260112000000_otel_traces.down.sql new file mode 100644 index 0000000000..3b2a767774 --- /dev/null +++ b/backend/migrations/20260112000000_otel_traces.down.sql @@ -0,0 +1,3 @@ +DROP INDEX IF EXISTS otel_traces_time_idx; +DROP INDEX IF EXISTS otel_traces_trace_time_idx; +DROP TABLE IF EXISTS otel_traces; diff --git a/backend/migrations/20260112000000_otel_traces.up.sql b/backend/migrations/20260112000000_otel_traces.up.sql new file mode 100644 index 0000000000..3da851f32a --- /dev/null +++ b/backend/migrations/20260112000000_otel_traces.up.sql @@ -0,0 +1,36 @@ +-- OpenTelemetry Span storage (all fields from proto::Span). +-- See: https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/trace/v1/trace.proto + +CREATE TABLE IF NOT EXISTS otel_traces ( + -- Identity fields (BYTEA for efficient storage and querying) + trace_id BYTEA NOT NULL, -- 16 bytes (proto: bytes) + span_id BYTEA NOT NULL, -- 8 bytes (proto: bytes) + trace_state TEXT NOT NULL DEFAULT '', -- W3C trace-context (proto: string) + parent_span_id BYTEA NOT NULL DEFAULT '', -- 8 bytes, empty if root span (proto: bytes) + flags INTEGER NOT NULL DEFAULT 0, -- W3C trace flags (proto: fixed32) + -- Core fields + name TEXT NOT NULL, -- operation name (proto: string) + kind INTEGER NOT NULL, -- SpanKind enum (proto: int32) + start_time_unix_nano BIGINT NOT NULL, -- (proto: fixed64, postgres has no u64) + end_time_unix_nano BIGINT NOT NULL, -- (proto: fixed64, postgres has no u64) + -- Attributes + attributes JSONB NOT NULL DEFAULT '[]', -- (proto: repeated KeyValue) + dropped_attributes_count INTEGER NOT NULL DEFAULT 0, -- (proto: uint32) + -- Events + events JSONB NOT NULL DEFAULT '[]', -- (proto: repeated Event) + dropped_events_count INTEGER NOT NULL DEFAULT 0, -- (proto: uint32) + -- Links + links JSONB NOT NULL DEFAULT '[]', -- (proto: repeated Link) + dropped_links_count INTEGER NOT NULL DEFAULT 0, -- (proto: uint32) + -- Status + status JSONB, -- (proto: optional Status message) + PRIMARY KEY (trace_id, span_id) +); + +-- Query spans by trace_id, ordered by time +CREATE INDEX IF NOT EXISTS otel_traces_trace_time_idx ON otel_traces (trace_id, start_time_unix_nano); + +-- Time-based cleanup (retention policy) +CREATE INDEX IF NOT EXISTS otel_traces_time_idx ON otel_traces (start_time_unix_nano); + +-- trace_id = job_id.as_bytes() diff --git a/backend/src/main.rs b/backend/src/main.rs index 187f04c8e1..065162211f 100644 --- a/backend/src/main.rs +++ b/backend/src/main.rs @@ -44,7 +44,7 @@ use windmill_common::{ EXPOSE_DEBUG_METRICS_SETTING, EXPOSE_METRICS_SETTING, EXTRA_PIP_INDEX_URL_SETTING, HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING, INDEXER_SETTING, INSTANCE_PYTHON_VERSION_SETTING, JOB_DEFAULT_TIMEOUT_SECS_SETTING, JWT_SECRET_SETTING, - KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MAVEN_REPOS_SETTING, + KEEP_JOB_DIR_SETTING, LICENSE_KEY_SETTING, MAVEN_REPOS_SETTING, OTEL_TRACING_PROXY_SETTING, MONITOR_LOGS_ON_OBJECT_STORE_SETTING, NO_DEFAULT_MAVEN_SETTING, NPM_CONFIG_REGISTRY_SETTING, NUGET_CONFIG_SETTING, OAUTH_SETTING, OTEL_SETTING, PIP_INDEX_URL_SETTING, POWERSHELL_REPO_PAT_SETTING, POWERSHELL_REPO_URL_SETTING, @@ -99,9 +99,9 @@ use crate::monitor::{ reload_bunfig_install_scopes_setting, reload_critical_alert_mute_ui_setting, reload_critical_error_channels_setting, reload_extra_pip_index_url_setting, reload_hub_api_secret_setting, reload_hub_base_url_setting, reload_job_default_timeout_setting, - reload_jwt_secret_setting, reload_license_key, reload_npm_config_registry_setting, - reload_pip_index_url_setting, reload_retention_period_setting, reload_scim_token_setting, - reload_smtp_config, reload_worker_config, MonitorIteration, + reload_jwt_secret_setting, reload_license_key, reload_otel_tracing_proxy_setting, + reload_npm_config_registry_setting, reload_pip_index_url_setting, reload_retention_period_setting, + reload_scim_token_setting, reload_smtp_config, reload_worker_config, MonitorIteration, }; #[cfg(feature = "parquet")] @@ -452,6 +452,7 @@ async fn windmill_main() -> anyhow::Result<()> { .unwrap_or(DEFAULT_NUM_WORKERS as i32) }; + // TODO: maybe gate behind debug_assertions? if num_workers > 1 && !std::env::var("WORKER_GROUP").is_ok_and(|x| x == "native") { println!( "We STRONGLY recommend using at most 1 worker per container, use at your own risks" @@ -802,6 +803,10 @@ Windmill Community Edition {GIT_VERSION} #[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() { @@ -1156,6 +1161,13 @@ Windmill Community Edition {GIT_VERSION} 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"); + send_delayed_killpill(&tx, 4, "OTEL tracing proxy setting change").await; + } + }, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING => { load_require_preexisting_user(&db).await; }, @@ -1364,6 +1376,42 @@ Windmill Community Edition {GIT_VERSION} Ok(()) as anyhow::Result<()> }; + let otel_tracing_proxy_f = async { + #[cfg(all(feature = "private", feature = "enterprise"))] + { + // Start OTEL tracing proxy for HTTP request interception + // Only enabled when: setting is on, worker mode (not server), and single worker (to avoid race conditions) + if worker_mode + && num_workers == 1 + && windmill_worker::OTEL_TRACING_PROXY_SETTINGS + .read() + .await + .enabled + { + if let Some(db) = conn.as_sql() { + tracing::info!( + "Starting OTEL tracing proxy (port will be dynamically assigned)" + ); + if let Err(e) = + windmill_worker::start_otel_tracing_proxy(db.clone(), otel_killpill_rx) + .await + { + tracing::error!("OTEL tracing proxy error: {}", e); + } + } + } else if windmill_worker::OTEL_TRACING_PROXY_SETTINGS + .read() + .await + .enabled + && num_workers > 1 + { + tracing::warn!("OTEL tracing proxy is enabled but num_workers > 1. Disabling to avoid race conditions. Set NUM_WORKERS=1 to enable."); + } + } + + Ok(()) as anyhow::Result<()> + }; + if server_mode { if let Some(db) = conn.as_sql() { schedule_stats(&db, &HTTP_CLIENT).await; @@ -1378,6 +1426,7 @@ Windmill Community Edition {GIT_VERSION} monitor_f, server_f, metrics_f, + otel_tracing_proxy_f, indexer_f, log_indexer_f )?; diff --git a/backend/src/monitor.rs b/backend/src/monitor.rs index 115dcf6d27..055feee0c2 100644 --- a/backend/src/monitor.rs +++ b/backend/src/monitor.rs @@ -53,7 +53,7 @@ use windmill_common::{ HUB_API_SECRET_SETTING, HUB_BASE_URL_SETTING, INSTANCE_PYTHON_VERSION_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, - NUGET_CONFIG_SETTING, OTEL_SETTING, PIP_INDEX_URL_SETTING, POWERSHELL_REPO_PAT_SETTING, + OTEL_TRACING_PROXY_SETTING, NUGET_CONFIG_SETTING, OTEL_SETTING, PIP_INDEX_URL_SETTING, POWERSHELL_REPO_PAT_SETTING, POWERSHELL_REPO_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, @@ -81,10 +81,10 @@ use windmill_common::{ use windmill_common::{client::AuthedClient, global_settings::APP_WORKSPACED_ROUTE_SETTING}; use windmill_queue::{cancel_job, get_queued_job_v2, SameWorkerPayload}; use windmill_worker::{ - handle_job_error, JobCompletedSender, SameWorkerSender, BUNFIG_INSTALL_SCOPES, - INSTANCE_PYTHON_VERSION, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, MAVEN_REPOS, NO_DEFAULT_MAVEN, - NPM_CONFIG_REGISTRY, NUGET_CONFIG, PIP_EXTRA_INDEX_URL, PIP_INDEX_URL, POWERSHELL_REPO_PAT, - POWERSHELL_REPO_URL, + handle_job_error, JobCompletedSender, OtelTracingProxySettings, SameWorkerSender, BUNFIG_INSTALL_SCOPES, + INSTANCE_PYTHON_VERSION, JOB_DEFAULT_TIMEOUT, KEEP_JOB_DIR, MAVEN_REPOS, + OTEL_TRACING_PROXY_SETTINGS, NO_DEFAULT_MAVEN, NPM_CONFIG_REGISTRY, NUGET_CONFIG, PIP_EXTRA_INDEX_URL, + PIP_INDEX_URL, POWERSHELL_REPO_PAT, POWERSHELL_REPO_URL, }; #[cfg(feature = "parquet")] @@ -320,6 +320,7 @@ pub async fn initial_load( reload_maven_repos_setting(&conn).await; reload_no_default_maven_setting(&conn).await; reload_ruby_repos_setting(&conn).await; + reload_otel_tracing_proxy_setting(&conn).await; } } @@ -778,6 +779,35 @@ pub async fn load_keep_job_dir(conn: &Connection) { }; } +pub async fn reload_otel_tracing_proxy_setting(conn: &Connection) { + match load_value_from_global_settings_with_conn(conn, OTEL_TRACING_PROXY_SETTING, true).await { + Ok(Some(settings)) => { + match serde_json::from_value::(settings) { + Ok(new_settings) => { + let mut current = OTEL_TRACING_PROXY_SETTINGS.write().await; + if current.enabled != new_settings.enabled + || current.enabled_languages != new_settings.enabled_languages + { + tracing::info!( + "OTEL tracing proxy settings changed: enabled={}, languages={:?}", + new_settings.enabled, + new_settings.enabled_languages + ); + *current = new_settings; + } + } + Err(e) => { + tracing::error!("Error parsing OTEL tracing proxy settings: {e:#}"); + } + } + } + Err(e) => { + tracing::error!("Error loading OTEL tracing proxy setting: {e:#}"); + } + _ => (), + }; +} + pub async fn load_require_preexisting_user(db: &DB) { let value = load_value_from_global_settings(db, REQUIRE_PREEXISTING_USER_FOR_OAUTH_SETTING).await; diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index dc23fab073..538ef8ce8c 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -9344,6 +9344,25 @@ paths: items: type: string + /w/{workspace}/jobs/get_otel_traces/{id}: + get: + summary: get OpenTelemetry traces for a job + operationId: getJobOtelTraces + tags: + - job + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/JobId" + responses: + "200": + description: list of OTEL Span objects (compatible with OpenTelemetry Span proto) + content: + application/json: + schema: + type: array + items: + type: object + /w/{workspace}/trigger/{trigger_kind}/resume_suspended_trigger_jobs/{trigger_path}: post: summary: resume all suspended jobs for a specific trigger @@ -17564,6 +17583,7 @@ components: - visible_to_owner - tag + ExportableCompletedJob: type: object description: Completed job with full data for export/import operations diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index daa13ca4b3..144963b12b 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -341,6 +341,7 @@ pub fn workspaced_service() -> Router { "/send_email_with_instance_smtp", post(send_email_with_instance_smtp), ) + .route("/get_otel_traces/:id", get(get_otel_traces)) } pub fn workspace_unauthed_service() -> Router { @@ -8835,3 +8836,63 @@ async fn delete_completed_job<'a>( ) .await; } + +async fn get_otel_traces( + OptAuthed(opt_authed): OptAuthed, + Extension(db): Extension, + Path((w_id, id)): Path<(String, Uuid)>, +) -> error::Result>> { + // Check job exists and user has permission to view it + let job = sqlx::query_scalar!( + "SELECT created_by FROM v2_job WHERE id = $1 AND workspace_id = $2", + id, + w_id + ) + .fetch_optional(&db) + .await?; + + match job { + Some(created_by) => { + if opt_authed.is_none() && created_by != "anonymous" { + return Err(Error::BadRequest( + "As a non logged in user, you can only see jobs ran by anonymous users" + .to_string(), + )); + } + } + None => { + return Err(Error::NotFound(format!("Job {} not found", id))); + } + } + + let trace_id = id.as_bytes().as_slice(); + + let traces = sqlx::query_scalar!( + r#"SELECT json_build_object( + 'trace_id', encode(trace_id, 'hex'), -- BYTEA to hex string + 'span_id', encode(span_id, 'hex'), -- BYTEA to hex string + 'parent_span_id', encode(parent_span_id, 'hex'), -- BYTEA to hex string + 'trace_state', trace_state, + 'flags', flags, + 'name', name, + 'kind', kind, + 'start_time_unix_nano', start_time_unix_nano, + 'end_time_unix_nano', end_time_unix_nano, + 'attributes', attributes, + 'dropped_attributes_count', dropped_attributes_count, + 'events', events, + 'dropped_events_count', dropped_events_count, + 'links', links, + 'dropped_links_count', dropped_links_count, + 'status', status + ) as "span!" + FROM otel_traces + WHERE trace_id = $1 + ORDER BY start_time_unix_nano"#, + trace_id + ) + .fetch_all(&db) + .await?; + + Ok(Json(traces)) +} diff --git a/backend/windmill-common/Cargo.toml b/backend/windmill-common/Cargo.toml index 1acf9f4c6a..314775c9b5 100644 --- a/backend/windmill-common/Cargo.toml +++ b/backend/windmill-common/Cargo.toml @@ -15,7 +15,7 @@ benchmark = [] parquet = ["dep:object_store", "dep:aws-sdk-sts", "dep:aws-smithy-types-convert", "dep:datafusion"] aws_auth = ["dep:aws-sdk-sts"] otel = ["dep:opentelemetry-semantic-conventions", "dep:opentelemetry-otlp", "dep:opentelemetry_sdk", - "dep:opentelemetry", "dep:tracing-opentelemetry", "dep:opentelemetry-appender-tracing", "dep:tonic"] + "dep:tracing-opentelemetry", "dep:opentelemetry-appender-tracing", "dep:tonic", "dep:opentelemetry"] smtp = ["dep:mail-send"] scoped_cache = [] cloud = [] diff --git a/backend/windmill-common/src/global_settings.rs b/backend/windmill-common/src/global_settings.rs index c5428e5c5c..890fba7821 100644 --- a/backend/windmill-common/src/global_settings.rs +++ b/backend/windmill-common/src/global_settings.rs @@ -46,6 +46,7 @@ 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 OTEL_TRACING_PROXY_SETTING: &str = "otel_tracing_proxy"; pub const APP_WORKSPACED_ROUTE_SETTING: &str = "app_workspaced_route"; pub const SECRET_BACKEND_SETTING: &str = "secret_backend"; diff --git a/backend/windmill-worker/Cargo.toml b/backend/windmill-worker/Cargo.toml index 60b4838b76..e5c41d7593 100644 --- a/backend/windmill-worker/Cargo.toml +++ b/backend/windmill-worker/Cargo.toml @@ -24,7 +24,7 @@ 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", "dep:deno_permissions", "dep:deno_io", "dep:deno_runtime", "dep:deno_telemetry", "dep:deno_error", "dep:winapi"] libffi_mac = ["dep:libffi-sys"] -otel = ["windmill-common/otel", "dep:opentelemetry"] +otel = ["windmill-common/otel", "dep:opentelemetry", "dep:tracing-opentelemetry"] dind = ["dep:bollard"] php = ["dep:windmill-parser-php"] mysql = ["dep:mysql_async"] @@ -76,6 +76,7 @@ futures.workspace = true async-recursion.workspace = true async-trait.workspace = true anyhow.workspace = true +derive_more.workspace = true itertools.workspace = true regex.workspace = true prometheus = { workspace = true, optional = true } @@ -136,9 +137,16 @@ process-wrap.workspace = true async-once-cell.workspace = true libloading = { workspace = true, optional = true } +opentelemetry-proto.workspace = true opentelemetry = { workspace = true, optional = true } +tracing-opentelemetry = { workspace = true, optional = true } bollard = { workspace = true, optional = true } oracle = { workspace = true, optional = true } +hudsucker.workspace = true +hyper-http-proxy.workspace = true +hyper-tls.workspace = true +hyper-util.workspace = true +rcgen.workspace = true [build-dependencies] deno_fetch = { workspace = true, optional = true } diff --git a/backend/windmill-worker/nsjail/download.py.config.proto b/backend/windmill-worker/nsjail/download.py.config.proto index c582642986..15fa807f66 100644 --- a/backend/windmill-worker/nsjail/download.py.config.proto +++ b/backend/windmill-worker/nsjail/download.py.config.proto @@ -100,3 +100,12 @@ exec_bin { arg: "/download_deps.sh" } +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} + diff --git a/backend/windmill-worker/nsjail/download.ruby.config.proto b/backend/windmill-worker/nsjail/download.ruby.config.proto index 7b9a04f912..e967623893 100644 --- a/backend/windmill-worker/nsjail/download.ruby.config.proto +++ b/backend/windmill-worker/nsjail/download.ruby.config.proto @@ -83,4 +83,11 @@ mount { is_bind: true } -{DEV} +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/download.rust.config.proto b/backend/windmill-worker/nsjail/download.rust.config.proto index a7e480eb2d..dfbc58f6a5 100644 --- a/backend/windmill-worker/nsjail/download.rust.config.proto +++ b/backend/windmill-worker/nsjail/download.rust.config.proto @@ -115,4 +115,11 @@ mount { rw: true } -{DEV} +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/lock.ruby.config.proto b/backend/windmill-worker/nsjail/lock.ruby.config.proto index 4d6035a772..86d46b8303 100644 --- a/backend/windmill-worker/nsjail/lock.ruby.config.proto +++ b/backend/windmill-worker/nsjail/lock.ruby.config.proto @@ -76,4 +76,11 @@ mount { is_bind: true } -{DEV} +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.bash.config.proto b/backend/windmill-worker/nsjail/run.bash.config.proto index 832579d950..8256f66f5b 100644 --- a/backend/windmill-worker/nsjail/run.bash.config.proto +++ b/backend/windmill-worker/nsjail/run.bash.config.proto @@ -138,3 +138,12 @@ iface_no_lo: true {SHARED_MOUNT} envar: "HOME=/tmp" + +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.bun.config.proto b/backend/windmill-worker/nsjail/run.bun.config.proto index 5c3738bc71..8cbee3dec9 100644 --- a/backend/windmill-worker/nsjail/run.bun.config.proto +++ b/backend/windmill-worker/nsjail/run.bun.config.proto @@ -161,4 +161,11 @@ iface_no_lo: true envar: "HOME=/tmp/{LANG}" +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.csharp.config.proto b/backend/windmill-worker/nsjail/run.csharp.config.proto index 11e991940c..7cd078c92f 100644 --- a/backend/windmill-worker/nsjail/run.csharp.config.proto +++ b/backend/windmill-worker/nsjail/run.csharp.config.proto @@ -112,3 +112,12 @@ mount { } {SHARED_MOUNT} + +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.go.config.proto b/backend/windmill-worker/nsjail/run.go.config.proto index c03b5ee73b..bdd7ebbbb8 100644 --- a/backend/windmill-worker/nsjail/run.go.config.proto +++ b/backend/windmill-worker/nsjail/run.go.config.proto @@ -102,4 +102,11 @@ iface_no_lo: true envar: "GOPATH=/tmp/.cache/go" envar: "HOME=/tmp/go" +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.nu.config.proto b/backend/windmill-worker/nsjail/run.nu.config.proto index 52742cf222..3fe7dc0cba 100644 --- a/backend/windmill-worker/nsjail/run.nu.config.proto +++ b/backend/windmill-worker/nsjail/run.nu.config.proto @@ -105,3 +105,12 @@ iface_no_lo: true {SHARED_MOUNT} envar: "HOME=/tmp" + +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.python3.config.proto b/backend/windmill-worker/nsjail/run.python3.config.proto index 31246b1aca..91cedb591d 100644 --- a/backend/windmill-worker/nsjail/run.python3.config.proto +++ b/backend/windmill-worker/nsjail/run.python3.config.proto @@ -141,4 +141,11 @@ envar: "LD_LIBRARY_PATH=/usr/local/lib:$LD_LIBRARY_PATH" envar: "PYTHONPATH={ADDITIONAL_PYTHON_PATHS}" envar: "HOME=/tmp" +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.ruby.config.proto b/backend/windmill-worker/nsjail/run.ruby.config.proto index 539b87e147..ec53c55809 100644 --- a/backend/windmill-worker/nsjail/run.ruby.config.proto +++ b/backend/windmill-worker/nsjail/run.ruby.config.proto @@ -109,5 +109,11 @@ iface_no_lo: true {SHARED_DEPENDENCIES} +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} -{DEV} +#{DEV} diff --git a/backend/windmill-worker/nsjail/run.rust.config.proto b/backend/windmill-worker/nsjail/run.rust.config.proto index 475a99873f..28959f1c94 100644 --- a/backend/windmill-worker/nsjail/run.rust.config.proto +++ b/backend/windmill-worker/nsjail/run.rust.config.proto @@ -106,5 +106,12 @@ mount { {SHARED_MOUNT} -{DEV} +mount { + src: "{TRACING_PROXY_CA_CERT_PATH}" + dst: "{TRACING_PROXY_CA_CERT_PATH}" + is_bind: true + mandatory: false +} + +#{DEV} diff --git a/backend/windmill-worker/src/bash_executor.rs b/backend/windmill-worker/src/bash_executor.rs index 81cf920082..76e5a92cf6 100644 --- a/backend/windmill-worker/src/bash_executor.rs +++ b/backend/windmill-worker/src/bash_executor.rs @@ -35,13 +35,15 @@ use crate::handle_child::run_future_with_polling_update_job_poller; use crate::{ common::{ - build_args_map, build_command_with_isolation, get_reserved_variables, read_file, read_file_content, start_child_process, - OccupancyMetrics, + build_args_map, build_command_with_isolation, get_reserved_variables, read_file, + read_file_content, start_child_process, OccupancyMetrics, DEV_CONF_NSJAIL, }, + get_proxy_envs_for_lang, handle_child::handle_child, - DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, + DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, TRACING_PROXY_CA_CERT_PATH, }; use windmill_common::client::AuthedClient; +use windmill_common::scripts::ScriptLang; lazy_static::lazy_static! { @@ -165,8 +167,7 @@ exit $exit_status .runnable_path .as_ref() .map(|x| { - !x.starts_with(INIT_SCRIPT_PATH_PREFIX) - && !x.starts_with(PERIODIC_SCRIPT_PATH_PREFIX) + !x.starts_with(INIT_SCRIPT_PATH_PREFIX) && !x.starts_with(PERIODIC_SCRIPT_PATH_PREFIX) }) .unwrap_or(true); @@ -178,7 +179,9 @@ exit $exit_status &NSJAIL_CONFIG_RUN_BASH_CONTENT .replace("{JOB_DIR}", job_dir) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{SHARED_MOUNT}", shared_mount), + .replace("{SHARED_MOUNT}", shared_mount) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; let mut cmd_args = vec![ "--config", @@ -193,7 +196,7 @@ exit $exit_status .current_dir(job_dir) .env_clear() .envs(reserved_variables) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Bash).await?) .env("PATH", PATH_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) .args(cmd_args) @@ -219,6 +222,7 @@ exit $exit_status .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Bash).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 404e221486..cf77467895 100644 --- a/backend/windmill-worker/src/bun_executor.rs +++ b/backend/windmill-worker/src/bun_executor.rs @@ -17,16 +17,17 @@ use crate::{ build_command_with_isolation, create_args_and_out_file, get_reserved_variables, parse_npm_config, read_file, read_file_content, read_result, start_child_process, write_file_binary, MaybeLock, OccupancyMetrics, StreamNotifier, + DEV_CONF_NSJAIL, }, handle_child::handle_child, BUNFIG_INSTALL_SCOPES, BUN_BUNDLE_CACHE_DIR, BUN_CACHE_DIR, BUN_NO_CACHE, BUN_PATH, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NODE_BIN_PATH, NODE_PATH, NPM_CONFIG_REGISTRY, - NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TZ_ENV, + NPM_PATH, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, get_proxy_envs_for_lang, }; use windmill_common::{ client::AuthedClient, s3_helpers::BundleFormat, - scripts::{id_to_codebase_info, CodebaseInfo}, + scripts::{id_to_codebase_info, CodebaseInfo, ScriptLang}, utils::WarnAfterExt, workspace_dependencies::WorkspaceDependenciesPrefetched, }; @@ -41,7 +42,6 @@ use tokio::io::AsyncReadExt; use windmill_common::{ error::{self, Result}, get_latest_hash_for_path, - scripts::ScriptLang, worker::{exists_in_cache, save_cache, write_file, Connection, DISABLE_BUNDLING}, DB, }; @@ -1375,7 +1375,9 @@ try {{ "/tmp/bun/shared" }, ), - ), + ) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); @@ -1416,6 +1418,7 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) .envs(common_bun_proc_envs) .env("PATH", PATH_ENV.as_str()) .args(args) @@ -1433,6 +1436,7 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) .envs(common_bun_proc_envs) .stdin(Stdio::null()) .stdout(Stdio::piped()) @@ -1463,6 +1467,7 @@ try {{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Bun).await?) .envs(common_bun_proc_envs) .stdin(Stdio::null()) .stdout(Stdio::piped()) diff --git a/backend/windmill-worker/src/common.rs b/backend/windmill-worker/src/common.rs index 95d6d7060e..923c86a454 100644 --- a/backend/windmill-worker/src/common.rs +++ b/backend/windmill-worker/src/common.rs @@ -50,6 +50,20 @@ use crate::agent_workers::UPDATE_PING_URL; use crate::{JOB_DEFAULT_TIMEOUT, MAX_RESULT_SIZE, MAX_TIMEOUT_DURATION, PATH_ENV}; use windmill_common::client::AuthedClient; +/// Additional nsjail config for development. Currently used for nix flake. +#[cfg(debug_assertions)] +pub const DEV_CONF_NSJAIL: &str = r#" +mount { + src: "/nix/store" + dst: "/nix/store" + is_bind: true + mandatory: false +} +"#; + +#[cfg(not(debug_assertions))] +pub const DEV_CONF_NSJAIL: &str = ""; + pub async fn build_args_map<'a>( job: &'a MiniPulledJob, client: &AuthedClient, diff --git a/backend/windmill-worker/src/csharp_executor.rs b/backend/windmill-worker/src/csharp_executor.rs index c2addba8e1..5b70ea9ddc 100644 --- a/backend/windmill-worker/src/csharp_executor.rs +++ b/backend/windmill-worker/src/csharp_executor.rs @@ -28,12 +28,14 @@ use windmill_queue::CanceledBy; use crate::{ common::{ build_command_with_isolation, check_executor_binary_exists, create_args_and_out_file, - get_reserved_variables, read_result, start_child_process, + get_reserved_variables, read_result, start_child_process, DEV_CONF_NSJAIL, }, - handle_child::handle_child, + handle_child::handle_child, get_proxy_envs_for_lang, CSHARP_CACHE_DIR, DISABLE_NSJAIL, DISABLE_NUSER, DOTNET_PATH, HOME_ENV, NSJAIL_PATH, - NUGET_CONFIG, PATH_ENV, TZ_ENV, + NUGET_CONFIG, PATH_ENV, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, }; +#[cfg(feature = "csharp")] +use windmill_common::scripts::ScriptLang; use crate::common::OccupancyMetrics; use windmill_common::client::AuthedClient; @@ -562,7 +564,9 @@ pub async fn handle_csharp_job( .replace("{CACHE_DIR}", CSHARP_CACHE_DIR) .replace("{CACHE_HASH}", &hash) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{SHARED_MOUNT}", shared_mount), + .replace("{SHARED_MOUNT}", shared_mount) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); nsjail_cmd @@ -570,6 +574,7 @@ pub async fn handle_csharp_job( .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::CSharp).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -600,6 +605,7 @@ pub async fn handle_csharp_job( .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::CSharp).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("DOTNET_CLI_TELEMETRY_OPTOUT", "true") diff --git a/backend/windmill-worker/src/deno_executor.rs b/backend/windmill-worker/src/deno_executor.rs index e80eec265d..23977e47a2 100644 --- a/backend/windmill-worker/src/deno_executor.rs +++ b/backend/windmill-worker/src/deno_executor.rs @@ -11,12 +11,12 @@ use crate::{ start_child_process, OccupancyMetrics, StreamNotifier, }, handle_child::handle_child, - DENO_CACHE_DIR, DENO_PATH, DISABLE_NSJAIL, HOME_ENV, NPM_CONFIG_REGISTRY, PATH_ENV, TZ_ENV, + get_proxy_envs_for_lang, DENO_CACHE_DIR, DENO_PATH, DISABLE_NSJAIL, HOME_ENV, NPM_CONFIG_REGISTRY, PATH_ENV, TZ_ENV, }; use windmill_common::client::AuthedClient; use tokio::{fs::File, io::AsyncReadExt, process::Command}; -use windmill_common::{error::Result, worker::write_file, BASE_URL}; +use windmill_common::{error::Result, scripts::ScriptLang, worker::write_file, BASE_URL}; use windmill_common::{ error::{self}, worker::Connection, @@ -92,6 +92,12 @@ async fn get_common_deno_proc_envs( crate::USERPROFILE_ENV.to_string(), ); } + + // 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); + } + return deno_envs; } diff --git a/backend/windmill-worker/src/go_executor.rs b/backend/windmill-worker/src/go_executor.rs index dc2b8d157a..b0475ee14a 100644 --- a/backend/windmill-worker/src/go_executor.rs +++ b/backend/windmill-worker/src/go_executor.rs @@ -1,4 +1,5 @@ -use crate::{common::MaybeLock, PROXY_ENVS}; +use crate::{common::MaybeLock, get_proxy_envs_for_lang}; +use windmill_common::scripts::ScriptLang; use std::{collections::HashMap, fs::DirBuilder, process::Stdio}; use itertools::Itertools; @@ -20,11 +21,11 @@ use windmill_queue::{append_logs, CanceledBy, MiniPulledJob}; use crate::{ common::{ build_command_with_isolation, capitalize, create_args_and_out_file, get_reserved_variables, - read_result, start_child_process, OccupancyMetrics, + read_result, start_child_process, OccupancyMetrics, DEV_CONF_NSJAIL, }, handle_child::handle_child, DISABLE_NSJAIL, DISABLE_NUSER, GOPRIVATE, GOPROXY, GO_BIN_CACHE_DIR, GO_CACHE_DIR, HOME_ENV, - NSJAIL_PATH, PATH_ENV, TZ_ENV, + NSJAIL_PATH, PATH_ENV, PROXY_ENVS, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, }; use windmill_common::client::AuthedClient; @@ -344,7 +345,9 @@ func Run(req Req) (interface{{}}, error){{ &NSJAIL_CONFIG_RUN_GO_CONTENT .replace("{JOB_DIR}", job_dir) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{SHARED_MOUNT}", shared_mount), + .replace("{SHARED_MOUNT}", shared_mount) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); nsjail_cmd @@ -352,6 +355,7 @@ func Run(req Req) (interface{{}}, error){{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Go).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -372,6 +376,7 @@ func Run(req Req) (interface{{}}, error){{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Go).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/handle_child.rs b/backend/windmill-worker/src/handle_child.rs index 8de0899c04..f54c673383 100644 --- a/backend/windmill-worker/src/handle_child.rs +++ b/backend/windmill-worker/src/handle_child.rs @@ -120,6 +120,11 @@ pub async fn handle_child( ) -> error::Result { let start = Instant::now(); + #[cfg(all(feature = "private", feature = "enterprise"))] + if crate::OTEL_TRACING_PROXY_SETTINGS.read().await.enabled { + crate::otel_tracing_proxy_ee::set_current_job_context(*job_id).await; + } + let pid = child.id(); #[cfg(target_os = "linux")] if let Some(pid) = pid { diff --git a/backend/windmill-worker/src/java_executor.rs b/backend/windmill-worker/src/java_executor.rs index b553b9eff1..6232f99be2 100644 --- a/backend/windmill-worker/src/java_executor.rs +++ b/backend/windmill-worker/src/java_executor.rs @@ -608,7 +608,6 @@ async fn run<'a>( .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) - .envs(PROXY_ENVS.clone()) .args(vec![ "--config", "run.config.proto", diff --git a/backend/windmill-worker/src/lib.rs b/backend/windmill-worker/src/lib.rs index dc82ddf2b5..90aa00bc08 100644 --- a/backend/windmill-worker/src/lib.rs +++ b/backend/windmill-worker/src/lib.rs @@ -34,6 +34,9 @@ mod global_cache; mod go_executor; mod graphql_executor; mod handle_child; +#[cfg(all(feature = "private", feature = "enterprise"))] +mod otel_tracing_proxy_ee; +mod otel_tracing_proxy_oss; pub mod job_logger; #[cfg(feature = "private")] pub mod job_logger_ee; @@ -78,6 +81,9 @@ pub use worker::*; pub use worker_lockfiles::{ process_relative_imports, trigger_dependents_to_recompute_dependencies, }; +pub use otel_tracing_proxy_oss::start_otel_tracing_proxy; +#[cfg(all(feature = "private", feature = "enterprise"))] +pub use otel_tracing_proxy_oss::{set_current_job_context, TRACING_PROXY_PORT}; pub use result_processor::handle_job_error; diff --git a/backend/windmill-worker/src/nu_executor.rs b/backend/windmill-worker/src/nu_executor.rs index 917c5333ee..7e17678426 100644 --- a/backend/windmill-worker/src/nu_executor.rs +++ b/backend/windmill-worker/src/nu_executor.rs @@ -14,11 +14,12 @@ use windmill_queue::{append_logs, CanceledBy, MiniPulledJob}; use crate::{ common::{ build_command_with_isolation, create_args_and_out_file, get_reserved_variables, - read_result, start_child_process, OccupancyMetrics, + read_result, start_child_process, OccupancyMetrics, DEV_CONF_NSJAIL, }, - handle_child, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, - PROXY_ENVS, + handle_child, get_proxy_envs_for_lang, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, + TRACING_PROXY_CA_CERT_PATH, }; +use windmill_common::scripts::ScriptLang; use windmill_common::client::AuthedClient; const NSJAIL_CONFIG_RUN_NU_CONTENT: &str = include_str!("../nsjail/run.nu.config.proto"); @@ -251,7 +252,9 @@ async fn run<'a>( .replace("{JOB_DIR}", job_dir) .replace("{NU_PATH}", &NU_PATH) .replace("{SHARED_MOUNT}", &shared_mount) - .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()), + .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); nsjail_cmd @@ -261,7 +264,7 @@ async fn run<'a>( .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Nu).await?) .args(vec![ "--config", "run.config.proto", @@ -300,7 +303,7 @@ async fn run<'a>( .env("BASE_INTERNAL_URL", base_internal_url) .envs(envs) .envs(reserved_variables) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Nu).await?) // TODO(v1): // "--plugins", // &format!( diff --git a/backend/windmill-worker/src/otel_tracing_proxy_oss.rs b/backend/windmill-worker/src/otel_tracing_proxy_oss.rs new file mode 100644 index 0000000000..8a0e6ba03f --- /dev/null +++ b/backend/windmill-worker/src/otel_tracing_proxy_oss.rs @@ -0,0 +1,13 @@ +//! OSS stubs for OTEL tracing proxy (EE feature) + +#[cfg(all(feature = "private", feature = "enterprise"))] +pub use crate::otel_tracing_proxy_ee::*; + +/// Start the OTEL tracing proxy (no-op in OSS) +#[cfg(not(all(feature = "private", feature = "enterprise")))] +pub async fn start_otel_tracing_proxy( + _db: windmill_common::DB, + _killpill_rx: tokio::sync::broadcast::Receiver<()>, +) -> Result<(), Box> { + Ok(()) +} diff --git a/backend/windmill-worker/src/python_executor.rs b/backend/windmill-worker/src/python_executor.rs index a2919fcfce..df0697ce37 100644 --- a/backend/windmill-worker/src/python_executor.rs +++ b/backend/windmill-worker/src/python_executor.rs @@ -31,6 +31,7 @@ use windmill_common::{ self, Error::{self}, }, + scripts::ScriptLang, utils::calculate_hash, worker::{ copy_dir_recursively, pad_string, split_python_requirements, write_file, Connection, @@ -127,12 +128,12 @@ use windmill_common::s3_helpers::OBJECT_STORE_SETTINGS; use crate::{ common::{ build_command_with_isolation, create_args_and_out_file, get_reserved_variables, read_file, - read_result, start_child_process, OccupancyMetrics, StreamNotifier, + read_result, start_child_process, OccupancyMetrics, StreamNotifier, DEV_CONF_NSJAIL, }, handle_child::handle_child, worker_utils::ping_job_status, PyV, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PIP_EXTRA_INDEX_URL, - PIP_INDEX_URL, PROXY_ENVS, PY_INSTALL_DIR, TZ_ENV, UV_CACHE_DIR, + PIP_INDEX_URL, PROXY_ENVS, PY_INSTALL_DIR, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, UV_CACHE_DIR, get_proxy_envs_for_lang, }; use windmill_common::client::AuthedClient; @@ -803,7 +804,9 @@ mount {{ .replace( "{ADDITIONAL_PYTHON_PATHS}", additional_python_paths_folders.as_str(), - ), + ) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), )?; } else { reserved_variables.insert("PYTHONPATH".to_string(), additional_python_paths_folders); @@ -822,7 +825,7 @@ mount {{ .env_clear() // inject PYTHONPATH here - for some reason I had to do it in nsjail conf .envs(reserved_variables) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Python3).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -848,6 +851,7 @@ mount {{ .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Python3).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -1372,6 +1376,8 @@ async fn spawn_uv_install( .replace("{PY_INSTALL_DIR}", &PY_INSTALL_DIR) .replace("{TARGET_DIR}", &venv_p) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL) .as_str(), )?; diff --git a/backend/windmill-worker/src/ruby_executor.rs b/backend/windmill-worker/src/ruby_executor.rs index 5b12aa5189..f0762a6e0d 100644 --- a/backend/windmill-worker/src/ruby_executor.rs +++ b/backend/windmill-worker/src/ruby_executor.rs @@ -24,12 +24,14 @@ use windmill_queue::{append_logs, CanceledBy, MiniPulledJob}; use crate::{ common::{ build_command_with_isolation, create_args_and_out_file, get_reserved_variables, - read_result, start_child_process, OccupancyMetrics, + read_result, start_child_process, OccupancyMetrics, DEV_CONF_NSJAIL, }, - handle_child::{self}, + handle_child::{self}, get_proxy_envs_for_lang, universal_pkg_installer::{par_install_language_dependencies_seq, RequiredDependency}, DISABLE_NSJAIL, DISABLE_NUSER, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, RUBY_CACHE_DIR, RUBY_REPOS, + TRACING_PROXY_CA_CERT_PATH, }; +use windmill_common::scripts::ScriptLang; lazy_static::lazy_static! { static ref RUBY_CONCURRENT_DOWNLOADS: usize = std::env::var("RUBY_CONCURRENT_DOWNLOADS").ok().map(|flag| flag.parse().unwrap_or(20)).unwrap_or(20); static ref RUBY_PATH: String = std::env::var("RUBY_PATH").unwrap_or_else(|_| "/usr/bin/ruby".to_string()); @@ -49,19 +51,6 @@ const NSJAIL_CONFIG_DOWNLOAD_RUBY_CONTENT: &str = include_str!("../nsjail/download.ruby.config.proto"); const NSJAIL_CONFIG_LOCK_RUBY_CONTENT: &str = include_str!("../nsjail/lock.ruby.config.proto"); -#[cfg(debug_assertions)] -const DEV_CONF_NSJAIL: &'static str = r#" -# Mount nix store for nixos to work properly -mount { - src: "/nix/store" - dst: "/nix/store" - is_bind: true - mandatory: false -} -"#; - -#[cfg(not(debug_assertions))] -const DEV_CONF_NSJAIL: &'static str = ""; #[allow(dead_code)] pub(crate) struct JobHandlerInput<'a> { @@ -349,7 +338,8 @@ Your Gemfile syntax will continue to work as-is." &NSJAIL_CONFIG_LOCK_RUBY_CONTENT .replace("{JOB_DIR}", job_dir) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{DEV}", DEV_CONF_NSJAIL), // .replace("{BUILD}", &build_dir), + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), // .replace("{BUILD}", &build_dir), )?; let mut cmd = Command::new(NSJAIL_PATH.as_str()); cmd.args(vec!["--config", &nsjail_proto, "--", BUNDLE_PATH.as_str()]); @@ -624,7 +614,8 @@ async fn install<'a>( &NSJAIL_CONFIG_DOWNLOAD_RUBY_CONTENT .replace("{TARGET}", &dependency.path) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{DEV}", DEV_CONF_NSJAIL), // .replace("{BUILD}", &build_dir), + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL), // .replace("{BUILD}", &build_dir), )?; let mut cmd = Command::new(NSJAIL_PATH.as_str()); cmd.args(vec!["--config", &nsjail_proto, "--", GEM_PATH.as_str()]); @@ -791,7 +782,8 @@ mount {{ .replace("{JOB_DIR}", job_dir) .replace("{SHARED_MOUNT}", &shared_mount) .replace("{SHARED_DEPENDENCIES}", &shared_deps) - .replace("{DEV}", DEV_CONF_NSJAIL) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()), )?; let mut cmd = Command::new(NSJAIL_PATH.as_str()); @@ -803,7 +795,7 @@ mount {{ .envs(envs) .envs(reserved_variables) .envs(RUBY_PROXY_ENVS.clone()) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Ruby).await?) .args(vec![ "--config", "run.config.proto", @@ -842,7 +834,7 @@ mount {{ .env("BASE_INTERNAL_URL", base_internal_url) .envs(reserved_variables) .envs(RUBY_PROXY_ENVS.clone()) - .envs(PROXY_ENVS.clone()) + .envs(get_proxy_envs_for_lang(&ScriptLang::Ruby).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 c43a7d2d97..90695145f8 100644 --- a/backend/windmill-worker/src/rust_executor.rs +++ b/backend/windmill-worker/src/rust_executor.rs @@ -20,12 +20,13 @@ use windmill_queue::{append_logs, CanceledBy}; use crate::{ common::{ build_command_with_isolation, check_executor_binary_exists, create_args_and_out_file, get_reserved_variables, - read_result, start_child_process, OccupancyMetrics, + read_result, start_child_process, OccupancyMetrics, DEV_CONF_NSJAIL, }, handle_child::handle_child, - DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, PROXY_ENVS, RUST_CACHE_DIR, - TZ_ENV, + get_proxy_envs_for_lang, DISABLE_NSJAIL, DISABLE_NUSER, HOME_ENV, NSJAIL_PATH, PATH_ENV, + PROXY_ENVS, RUST_CACHE_DIR, TRACING_PROXY_CA_CERT_PATH, TZ_ENV, }; +use windmill_common::scripts::ScriptLang; use windmill_common::client::AuthedClient; #[cfg(windows)] @@ -52,19 +53,6 @@ lazy_static::lazy_static! { static ref RUSTUP_HOME_DEFAULT: String = format!("{}\\.rustup", *HOME_DIR); } -#[cfg(debug_assertions)] -const DEV_CONF_NSJAIL: &'static str = r#" -# Mount nix store for nixos to work properly -mount { - src: "/nix/store" - dst: "/nix/store" - is_bind: true - mandatory: false -} -"#; - -#[cfg(not(debug_assertions))] -const DEV_CONF_NSJAIL: &'static str = ""; #[cfg(not(windows))] lazy_static::lazy_static! { @@ -341,7 +329,8 @@ pub async fn build_rust_crate( .replace("{JOB_DIR}", job_dir) .replace("{CACHE_DIR}", RUST_CACHE_DIR) .replace("{CARGO_HOME}", CARGO_HOME.as_str()) - .replace("{DEV}", DEV_CONF_NSJAIL) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL) .replace("{BUILD}", &build_dir), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); @@ -549,7 +538,8 @@ pub async fn handle_rust_job( .replace("{CACHE_DIR}", RUST_CACHE_DIR) .replace("{CACHE_HASH}", &hash) .replace("{CLONE_NEWUSER}", &(!*DISABLE_NUSER).to_string()) - .replace("{DEV}", DEV_CONF_NSJAIL) + .replace("{TRACING_PROXY_CA_CERT_PATH}", TRACING_PROXY_CA_CERT_PATH) + .replace("#{DEV}", DEV_CONF_NSJAIL) .replace("{SHARED_MOUNT}", shared_mount), )?; let mut nsjail_cmd = Command::new(NSJAIL_PATH.as_str()); @@ -558,6 +548,7 @@ pub async fn handle_rust_job( .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Rust).await?) .env("PATH", PATH_ENV.as_str()) .env("TZ", TZ_ENV.as_str()) .env("BASE_INTERNAL_URL", base_internal_url) @@ -573,6 +564,7 @@ pub async fn handle_rust_job( .env_clear() .envs(envs) .envs(reserved_variables) + .envs(get_proxy_envs_for_lang(&ScriptLang::Rust).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 4cdb0ebf6d..5cc4972c32 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -55,7 +55,7 @@ use windmill_common::METRICS_ENABLED; use serde::{Deserialize, Serialize}; use sqlx::types::Json; use std::{ - collections::HashMap, + collections::{HashMap, HashSet}, fmt::Display, sync::{ atomic::{AtomicBool, AtomicU16, Ordering}, @@ -228,6 +228,9 @@ pub const GO_BIN_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "gobin"); pub const POWERSHELL_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "powershell"); pub const COMPOSER_CACHE_DIR: &str = concatcp!(ROOT_CACHE_DIR, "composer"); +pub const TRACING_PROXY_CA_CERT_PATH: &str = + concatcp!(ROOT_CACHE_NOMOUNT_DIR, "tracing_proxy_ca.pem"); + const NUM_SECS_PING: u64 = 5; const NUM_SECS_READINGS: u64 = 60; @@ -248,6 +251,16 @@ const VACUUM_PERIOD: u32 = 50000; pub const MAX_BUFFERED_DEDICATED_JOBS: usize = 3; +/// Per-language OTEL tracing proxy configuration. +/// Default languages are configured in frontend instanceSettings.ts +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct OtelTracingProxySettings { + #[serde(default)] + pub enabled: bool, + #[serde(default)] + pub enabled_languages: HashSet, +} + #[cfg(feature = "prometheus")] lazy_static::lazy_static! { @@ -505,6 +518,7 @@ lazy_static::lazy_static! { pub static ref HTTP_PROXY: Option = std::env::var("http_proxy").ok().or(std::env::var("HTTP_PROXY").ok()); pub static ref HTTPS_PROXY: Option = std::env::var("https_proxy").ok().or(std::env::var("HTTPS_PROXY").ok()); + /// Static proxy environment variables from env vars (for languages not using dynamic OTEL tracing proxy config) pub static ref PROXY_ENVS: Vec<(&'static str, String)> = { let mut proxy_env = Vec::new(); if let Some(no_proxy) = NO_PROXY.as_ref() { @@ -520,6 +534,9 @@ lazy_static::lazy_static! { } proxy_env }; + + /// Per-language OTEL tracing proxy settings (configured via instance settings) + pub static ref OTEL_TRACING_PROXY_SETTINGS: Arc> = Arc::new(RwLock::new(OtelTracingProxySettings::default())); pub static ref WHITELIST_ENVS: HashMap = { windmill_common::worker::load_env_vars( windmill_common::worker::load_whitelist_env_vars_from_env(), @@ -609,6 +626,50 @@ lazy_static::lazy_static! { type Envs = Vec<(String, String)>; +/// Check if OTEL tracing proxy is enabled for a specific language (EE only) +pub async fn is_otel_tracing_proxy_enabled_for_lang(lang: &ScriptLang) -> bool { + cfg!(all(feature = "private", feature = "enterprise")) && { + let settings = OTEL_TRACING_PROXY_SETTINGS.read().await; + settings.enabled && settings.enabled_languages.contains(lang) + } +} + +/// Get proxy environment variables for job execution for a specific language. +/// When OTEL tracing proxy is enabled for this language, routes all traffic through the proxy. +/// Otherwise, uses the standard HTTP_PROXY/HTTPS_PROXY from environment. +pub async fn get_proxy_envs_for_lang(lang: &ScriptLang) -> 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; + } + let _ = lang; + 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 + .read() + .await + .ok_or_else(|| anyhow::anyhow!("OTEL tracing proxy port not initialized"))?; + let proxy_url = format!("http://127.0.0.1:{}", port); + Ok(vec![ + ("HTTP_PROXY", proxy_url.clone()), + ("HTTPS_PROXY", proxy_url.clone()), + // Lowercase variants for Ruby and other runtimes that check lowercase first + ("http_proxy", proxy_url.clone()), + ("https_proxy", proxy_url), + ("NO_PROXY", "".to_string()), + ("no_proxy", "".to_string()), + // CA cert for various runtimes to trust the tracing proxy + ("SSL_CERT_FILE", TRACING_PROXY_CA_CERT_PATH.to_string()), + ("REQUESTS_CA_BUNDLE", TRACING_PROXY_CA_CERT_PATH.to_string()), + ("NODE_EXTRA_CA_CERTS", TRACING_PROXY_CA_CERT_PATH.to_string()), + ("CURL_CA_BUNDLE", TRACING_PROXY_CA_CERT_PATH.to_string()), + ("DENO_CERT", TRACING_PROXY_CA_CERT_PATH.to_string()), + ]) +} + #[cfg(windows)] lazy_static::lazy_static! { pub static ref SYSTEM_ROOT: String = std::env::var("SystemRoot").unwrap_or_else(|_| "C:\\Windows".to_string()); @@ -2689,8 +2750,6 @@ pub async fn handle_queued_job( flow_runners: Option>, #[cfg(feature = "benchmark")] _bench: &mut BenchmarkIter, ) -> windmill_common::error::Result { - // Extract the active span from the context - if job.canceled_by.is_some() { return Err(Error::JsonErr(canceled_job_to_result(&job))); } @@ -3136,7 +3195,7 @@ pub fn build_envs( }; for (k, v) in PROXY_ENVS.iter() { - envs.insert(k.to_string(), v.to_string()); + envs.insert(k.to_string(), v.clone()); } Ok(envs) diff --git a/flake.nix b/flake.nix index 35156c0ac8..75bf82de8c 100644 --- a/flake.nix +++ b/flake.nix @@ -287,6 +287,7 @@ JAVA_PATH = "${pkgs.jdk21}/bin/java"; JAVAC_PATH = "${pkgs.jdk21}/bin/javac"; COURSIER_PATH = "${coursier}/coursier"; + BASH_PATH = "bash"; RUBY_PATH = "${pkgs.ruby}/bin/ruby"; RUBY_BUNDLE_PATH = "${pkgs.ruby}/bin/bundle"; RUBY_GEM_PATH = "${pkgs.ruby}/bin/gem"; diff --git a/frontend/src/lib/components/FlowStatusViewerInner.svelte b/frontend/src/lib/components/FlowStatusViewerInner.svelte index 9680c13ce5..2e157cdcd7 100644 --- a/frontend/src/lib/components/FlowStatusViewerInner.svelte +++ b/frontend/src/lib/components/FlowStatusViewerInner.svelte @@ -48,6 +48,7 @@ } from './graph/renderers/nodes/AIToolNode.svelte' import JobAssetsViewer from './assets/JobAssetsViewer.svelte' import McpToolCallDetails from './McpToolCallDetails.svelte' + import JobOtelTraces from './JobOtelTraces.svelte' import { SelectionManager } from './graph/selectionUtils.svelte' import { useThrottle } from 'runed' @@ -103,7 +104,7 @@ refreshGlobal: (moduleId: string, clear: boolean, root: string) => Promise updateGlobalRefresh: (moduleId: string, updateFn: (clear, root) => Promise) => void job?: (Job & { result_stream?: string }) | undefined - rightColumnSelect?: 'timeline' | 'node_status' | 'node_definition' | 'user_states' + rightColumnSelect?: 'timeline' | 'node_status' | 'node_definition' | 'user_states' | 'tracing' localModuleStates?: Record localDurationStatuses?: Record onResultStreamUpdate?: ({ @@ -1756,7 +1757,7 @@ bind:expandedSubflows onSelect={(e) => { console.log('onSelect', e) - if (rightColumnSelect != 'node_definition') { + if (rightColumnSelect != 'node_definition' && rightColumnSelect != 'tracing') { rightColumnSelect = 'node_status' } if (typeof e == 'string') { @@ -1811,6 +1812,7 @@ {#if Object.keys(job?.flow_status?.user_states ?? {}).length > 0} {/if} + {#if rightColumnSelect == 'timeline'} + {:else if rightColumnSelect == 'tracing'} + {@const node = selectedNode ? localModuleStates[selectedNode] : undefined} + {#if node?.job_id} + + {:else} +
Select a node with a job to see HTTP request traces
+ {/if} {/if} diff --git a/frontend/src/lib/components/InstanceSetting.svelte b/frontend/src/lib/components/InstanceSetting.svelte index 42a77b88dd..470c187b91 100644 --- a/frontend/src/lib/components/InstanceSetting.svelte +++ b/frontend/src/lib/components/InstanceSetting.svelte @@ -3,6 +3,8 @@ import { enterpriseLicense, isCriticalAlertsUIOpen } from '$lib/stores' import { AlertCircle, BadgeCheck, BadgeX, Info } from 'lucide-svelte' import type { Setting } from './instanceSettings' + import { OTEL_TRACING_PROXY_LANGUAGES } from './instanceSettings' + import { LanguageIcon } from './common/languageIcons' import Tooltip from './Tooltip.svelte' import ObjectStoreConfigSettings from './ObjectStoreConfigSettings.svelte' import { sendUserToast } from '$lib/toast' @@ -831,6 +833,39 @@ --> {/if} + {:else if setting.fieldType == 'otel_tracing_proxy'} + {@const tracingProxyVal = $values[setting.key] ?? { enabled: false, enabled_languages: [...OTEL_TRACING_PROXY_LANGUAGES] }} +
+ { + $values[setting.key] = { ...tracingProxyVal, enabled: e.detail } + }} + options={{ right: 'Enabled' }} + /> + {#if tracingProxyVal.enabled} +
+ {#each OTEL_TRACING_PROXY_LANGUAGES as lang (lang)} + {@const isEnabled = (tracingProxyVal.enabled_languages ?? []).includes(lang)} + + {/each} +
+ {/if} +
{:else if setting.fieldType == 'object_store_config'}
diff --git a/frontend/src/lib/components/JobOtelTraces.svelte b/frontend/src/lib/components/JobOtelTraces.svelte new file mode 100644 index 0000000000..d57a7158e1 --- /dev/null +++ b/frontend/src/lib/components/JobOtelTraces.svelte @@ -0,0 +1,311 @@ + + +
+ {#if loading} + + {:else if error} + {error} + {:else if traces.length === 0} +
+ +

No HTTP requests captured

+

+ This job did not make any HTTP/HTTPS requests, or HTTP Request Tracing is not enabled in instance settings. +

+
+ {:else} +
+
+

Traces ({traces.length} spans)

+ +
+ + +
+
+
+
Operation
+
Status
+
Timeline
+
Duration
+
+
+ +
+ {#each traces as span (span.span_id)} + {@const durationNs = span.end_time_unix_nano - span.start_time_unix_nano} + {@const startOffset = + ((span.start_time_unix_nano - timelineMetrics.minTime) / + timelineMetrics.totalDuration) * + 100} + {@const width = (durationNs / timelineMetrics.totalDuration) * 100} + {@const statusCode = span.status?.code ?? 0} + +
+ + + {#if expandedSpans.has(span.span_id)} + {@const parsedAttrs = parseAttributes(span.attributes)} +
+
+
+

Trace ID

+

{span.trace_id}

+
+
+

Span ID

+

{span.span_id}

+
+ {#if span.parent_span_id} +
+

Parent Span ID

+

{span.parent_span_id}

+
+ {/if} +
+

Kind

+

{getKindLabel(span.kind)}

+
+
+

Start Time

+

{formatTimestamp(span.start_time_unix_nano)}

+
+
+

End Time

+

{formatTimestamp(span.end_time_unix_nano)}

+
+ {#if span.status?.message} +
+

Status Message

+

{span.status.message}

+
+ {/if} + {#if Object.keys(parsedAttrs).length > 0} +
+

Attributes

+
+ {#each Object.entries(parsedAttrs) as [key, value]} +
+ {key}: + {typeof value === 'object' ? JSON.stringify(value) : value} +
+ {/each} +
+
+ {/if} +
+
+ {/if} +
+ {/each} +
+
+
+ {/if} +
diff --git a/frontend/src/lib/components/instanceSettings.ts b/frontend/src/lib/components/instanceSettings.ts index cb8b38f95f..44d0d3b2d3 100644 --- a/frontend/src/lib/components/instanceSettings.ts +++ b/frontend/src/lib/components/instanceSettings.ts @@ -1,5 +1,18 @@ import type { ButtonType } from './common/button/model' +// Languages that support HTTP request tracing via OTEL proxy +export const OTEL_TRACING_PROXY_LANGUAGES = [ + 'python3', + 'deno', + 'bun', + 'go', + 'bash', + 'rust', + 'csharp', + 'nu', + 'ruby' +] as const + export interface Setting { label: string description?: string @@ -34,6 +47,7 @@ export interface Setting { | 'smtp_connect' | 'indexer_rates' | 'otel' + | 'otel_tracing_proxy' | 'secret_backend' storage: SettingStorage advancedToggle?: { @@ -446,7 +460,16 @@ export const settings: Record = { storage: 'setting', ee_only: '' }, - + { + label: 'HTTP Request Tracing', + description: + 'Capture HTTP/HTTPS requests from job scripts as OpenTelemetry spans. Visible in job details and exported to your OTEL collector if configured. Toggling restarts workers.', + key: 'otel_tracing_proxy', + fieldType: 'otel_tracing_proxy', + storage: 'setting', + ee_only: 'HTTP Request Tracing is an EE feature', + defaultValue: () => ({ enabled: false, enabled_languages: [...OTEL_TRACING_PROXY_LANGUAGES] }) + }, { label: 'Prometheus', description: diff --git a/frontend/src/lib/components/scriptEditor/LogPanel.svelte b/frontend/src/lib/components/scriptEditor/LogPanel.svelte index be9ddfcff0..013c87b6c9 100644 --- a/frontend/src/lib/components/scriptEditor/LogPanel.svelte +++ b/frontend/src/lib/components/scriptEditor/LogPanel.svelte @@ -25,6 +25,7 @@ import type Editor from '../Editor.svelte' import type DiffEditor from '../DiffEditor.svelte' import ScriptFix from '../copilot/ScriptFix.svelte' + import JobOtelTraces from '../JobOtelTraces.svelte' import Cell from '../table/Cell.svelte' import DataTable from '../table/DataTable.svelte' import Head from '../table/Head.svelte' @@ -127,6 +128,7 @@ {#if showCaptures && customUi?.disableTriggerCaptures !== true} {/if} + {#snippet content()}
@@ -298,6 +300,15 @@ {#if selectedTab === 'captures'} {@render capturesTab?.()} {/if} + {#if selectedTab === 'tracing'} + {#if previewJob?.id} + + {:else} +
+ Run a preview to see HTTP request traces +
+ {/if} + {/if}
{/snippet} diff --git a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte index c4fbbb3156..3452da8c9b 100644 --- a/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte +++ b/frontend/src/routes/(root)/(logged)/run/[...run]/+page.svelte @@ -91,13 +91,14 @@ import RunBadges from '$lib/components/runs/RunBadges.svelte' import { twMerge } from 'tailwind-merge' import FlowRestartButton from '$lib/components/FlowRestartButton.svelte' + import JobOtelTraces from '$lib/components/JobOtelTraces.svelte' let job: (Job & { result?: any; result_stream?: string }) | undefined = $state() let jobUpdateLastFetch: Date | undefined = $state() let scriptProgress: number | undefined = $state(undefined) let currentJobIsLongRunning: boolean = $state(false) - let viewTab: 'result' | 'logs' | 'code' | 'stats' | 'assets' = $state('result') + let viewTab: 'result' | 'logs' | 'code' | 'stats' | 'assets' | 'traces' = $state('result') let selectedJobStep: string | undefined = $state(undefined) let selectedJobStepIsTopLevel: boolean | undefined = $state(undefined) @@ -770,6 +771,7 @@ + {#if isScriptPreview(job?.job_kind)} @@ -799,6 +801,10 @@
+ {:else if viewTab == 'traces'} +
+ +
{:else if viewTab == 'code'} {#if job && 'raw_code' in job && job.raw_code}