From 9fd99c8cd3c71fae14216e8cc3234a3ec9b3b18d Mon Sep 17 00:00:00 2001 From: Kris Hicks Date: Mon, 17 Aug 2026 11:00:53 -0700 Subject: [PATCH] feat(podman): export driver traces over OTLP Mirror the VM driver tracing setup for Podman. Export standalone driver spans through OTLP/gRPC as the distinct openshell-driver-podman service, propagate W3C context across ComputeDriver RPCs, record bounded RPC names and failures, and flush buffered spans during graceful shutdown. Podman still runs in-process when selected as a built-in gateway driver. Add a temporary tracing shim that partitions gateway and Podman spans by target into separate tracer providers while preserving their shared trace and parentage. The shim also emits the same ComputeDriver server boundary that the tonic layer emits out of process, keeping the observable trace shape stable when Podman is eventually extracted. Trace container create preparation, image and storage setup, lifecycle operations, and cleanup. Document the service boundary and cover it with isolated and repeated tracing tests. Signed-off-by: Kris Hicks --- Cargo.lock | 7 + crates/openshell-driver-podman/Cargo.toml | 8 + crates/openshell-driver-podman/README.md | 6 + crates/openshell-driver-podman/src/driver.rs | 643 +++++++++++++----- crates/openshell-driver-podman/src/grpc.rs | 593 +++++++++++++--- crates/openshell-driver-podman/src/lib.rs | 1 + crates/openshell-driver-podman/src/main.rs | 97 ++- .../src/otel_tracing.rs | 299 ++++++++ .../openshell-driver-vm/src/otel_tracing.rs | 19 +- crates/openshell-otel/Cargo.toml | 1 + crates/openshell-otel/src/grpc.rs | 117 +++- crates/openshell-otel/src/lib.rs | 121 +++- crates/openshell-server/src/cli.rs | 9 +- crates/openshell-server/src/compute/mod.rs | 2 +- crates/openshell-server/src/lib.rs | 21 +- crates/openshell-server/src/otel_tracing.rs | 8 +- crates/openshell-server/src/tracing_setup.rs | 59 +- 17 files changed, 1752 insertions(+), 259 deletions(-) create mode 100644 crates/openshell-driver-podman/src/otel_tracing.rs diff --git a/Cargo.lock b/Cargo.lock index c30f890914..5c19d01dfd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3822,12 +3822,17 @@ version = "0.0.0" dependencies = [ "clap", "futures", + "http 1.4.0", "http-body-util", "hyper 1.9.0", "hyper-util", "miette", "nix 0.29.0", "openshell-core", + "openshell-otel", + "opentelemetry", + "opentelemetry-proto", + "opentelemetry_sdk", "prost-types", "rustix 1.1.4", "serde", @@ -3837,7 +3842,9 @@ dependencies = [ "tokio", "tokio-stream", "tonic", + "tower-http 0.6.8", "tracing", + "tracing-opentelemetry", "tracing-subscriber", "url", ] diff --git a/crates/openshell-driver-podman/Cargo.toml b/crates/openshell-driver-podman/Cargo.toml index c989c13945..7442979711 100644 --- a/crates/openshell-driver-podman/Cargo.toml +++ b/crates/openshell-driver-podman/Cargo.toml @@ -16,6 +16,7 @@ path = "src/main.rs" [dependencies] openshell-core = { path = "../openshell-core", default-features = false, features = ["driver-extraction"] } +openshell-otel = { path = "../openshell-otel" } tokio = { workspace = true } tonic = { workspace = true, features = ["transport"] } @@ -32,11 +33,18 @@ nix = { workspace = true } rustix = { workspace = true } tracing = { workspace = true } tracing-subscriber = { workspace = true } +tracing-opentelemetry = { workspace = true } +opentelemetry = { workspace = true } +opentelemetry_sdk = { workspace = true } +tower-http = { workspace = true } +http = { workspace = true } thiserror = { workspace = true } miette = { workspace = true } url = { workspace = true } [dev-dependencies] +opentelemetry-proto = { version = "0.32", default-features = false, features = ["gen-tonic", "trace"] } +opentelemetry_sdk = { workspace = true, features = ["testing"] } prost-types = { workspace = true } temp-env = "0.3" tokio = { workspace = true, features = ["test-util"] } diff --git a/crates/openshell-driver-podman/README.md b/crates/openshell-driver-podman/README.md index c7778e5ca1..a8acd62387 100644 --- a/crates/openshell-driver-podman/README.md +++ b/crates/openshell-driver-podman/README.md @@ -7,6 +7,12 @@ driver runs in-process within the gateway server and delegates all sandbox isolation enforcement to the `openshell-sandbox` supervisor binary, which is sideloaded into each container via an OCI image volume mount. +When the gateway configures `[openshell.gateway.otlp]`, Podman compute-driver +spans export to the same OTLP/gRPC collector with the service name +`openshell-driver-podman`. The driver preserves the gateway trace context and +uses the same compute-driver RPC span names in its in-process and standalone +forms. + Before creating the container, the driver inspects the final sandbox image and captures its immutable image ID and raw OCI `Config.User`. Container creation uses that image ID with pulling disabled, preventing a mutable tag from changing diff --git a/crates/openshell-driver-podman/src/driver.rs b/crates/openshell-driver-podman/src/driver.rs index 23175b3bf4..bb979b07b2 100644 --- a/crates/openshell-driver-podman/src/driver.rs +++ b/crates/openshell-driver-podman/src/driver.rs @@ -33,7 +33,7 @@ use std::net::{IpAddr, SocketAddr}; use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Duration; -use tracing::{debug, info, warn}; +use tracing::{Instrument as _, debug, info, warn}; use url::Url; const STOP_COMPLETION_POLL_INTERVAL: Duration = Duration::from_millis(50); @@ -682,7 +682,18 @@ impl PodmanComputeDriver { } /// Create a sandbox container. + #[tracing::instrument( + name = "podman.create_sandbox", + skip(self, sandbox), + fields( + otel.name = "podman.create_sandbox", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox.id, + sandbox.name = %sandbox.name, + ) + )] pub async fn create_sandbox(&self, sandbox: &DriverSandbox) -> Result<(), ComputeDriverError> { + let span_status = openshell_otel::ErrorStatusGuard::current(); if sandbox.name.is_empty() { return Err(ComputeDriverError::Precondition( "sandbox name is required".into(), @@ -709,84 +720,119 @@ impl PodmanComputeDriver { "Creating sandbox container" ); - // 1a. Pull the supervisor image if needed. The supervisor binary - // is shipped in a standalone OCI image and mounted into sandbox - // containers via Podman's type=image mount. Refresh mutable tags - // like latest/dev, but avoid registry checks for pinned images. - let supervisor_pull_policy = supervisor_image_pull_policy(&self.config.supervisor_image); - info!( - image = %self.config.supervisor_image, - policy = supervisor_pull_policy, - "Ensuring supervisor image" - ); - self.client - .pull_image(&self.config.supervisor_image, supervisor_pull_policy) - .await - .map_err(ComputeDriverError::from)?; - - // 1b. Pull the sandbox image if needed (Podman does not pull on create). - let image = container::resolve_image(sandbox, &self.config); - if image.is_empty() { - return Err(ComputeDriverError::Precondition( - "no sandbox image configured: set default_image in [openshell.drivers.podman] \ - or provide an image in the sandbox template" - .to_string(), - )); - } - let pull_policy = self.config.image_pull_policy.as_str(); - info!(image = %image, policy = %pull_policy, "Ensuring sandbox image"); - self.client - .pull_image(image, pull_policy) - .await - .map_err(ComputeDriverError::from)?; - let inspected_image = self - .client - .inspect_image(image) - .await - .map_err(ComputeDriverError::from)?; - if inspected_image.id.is_empty() { - return Err(ComputeDriverError::Precondition(format!( - "podman image '{image}' inspection did not return an immutable image ID" - ))); - } - let image_user = inspected_image - .config - .as_ref() - .map_or("", |config| config.user.as_str()); + let (image, immutable_image_id, image_user) = async { + let phase_status = openshell_otel::ErrorStatusGuard::current(); + let result = async { + // The supervisor binary is shipped in a standalone OCI image and + // mounted into sandbox containers via Podman's type=image mount. + let supervisor_pull_policy = + supervisor_image_pull_policy(&self.config.supervisor_image); + info!( + image = %self.config.supervisor_image, + policy = supervisor_pull_policy, + "Ensuring supervisor image" + ); + self.client + .pull_image(&self.config.supervisor_image, supervisor_pull_policy) + .await + .map_err(ComputeDriverError::from)?; - for image in - container::podman_driver_image_mount_sources(sandbox, self.config.enable_bind_mounts) + // Podman does not pull the sandbox image on container creation. + let image = container::resolve_image(sandbox, &self.config); + if image.is_empty() { + return Err(ComputeDriverError::Precondition( + "no sandbox image configured: set default_image in \ + [openshell.drivers.podman] or provide an image in the sandbox template" + .to_string(), + )); + } + let pull_policy = self.config.image_pull_policy.as_str(); + info!(image = %image, policy = %pull_policy, "Ensuring sandbox image"); + self.client + .pull_image(image, pull_policy) + .await + .map_err(ComputeDriverError::from)?; + let inspected_image = self + .client + .inspect_image(image) + .await + .map_err(ComputeDriverError::from)?; + if inspected_image.id.is_empty() { + return Err(ComputeDriverError::Precondition(format!( + "podman image '{image}' inspection did not return an immutable image ID" + ))); + } + let image_user = inspected_image + .config + .as_ref() + .map_or_else(String::new, |config| config.user.clone()); + + for mount_image in container::podman_driver_image_mount_sources( + sandbox, + self.config.enable_bind_mounts, + ) .map_err(ComputeDriverError::Precondition)? - { - info!(image = %image, policy = %pull_policy, "Ensuring image mount source"); - self.client - .pull_image(&image, pull_policy) - .await - .map_err(ComputeDriverError::from)?; - } + { + info!(image = %mount_image, policy = %pull_policy, "Ensuring image mount source"); + self.client + .pull_image(&mount_image, pull_policy) + .await + .map_err(ComputeDriverError::from)?; + } - // 2. Create workspace volume and per-sandbox token secret. - if let Err(e) = self.client.create_volume(&vol_name).await { - return Err(ComputeDriverError::from(e)); + Ok((image.to_string(), inspected_image.id, image_user)) + } + .await; + phase_status.finish(result) } - let token_secret_name = match create_sandbox_token_secret(&self.client, sandbox).await { - Ok(name) => name, - Err(e) => { - let _ = self.client.remove_volume(&vol_name).await; - return Err(e); + .instrument(tracing::info_span!( + "podman.prepare_images", + otel.name = "podman.prepare_images", + otel.status_code = tracing::field::Empty, + )) + .await?; + + // Create workspace volume and per-sandbox token secret. + let (token_secret_name, proxy_auth_secret_name) = async { + let phase_status = openshell_otel::ErrorStatusGuard::current(); + let result = async { + self.client + .create_volume(&vol_name) + .await + .map_err(ComputeDriverError::from)?; + let token_secret_name = + match create_sandbox_token_secret(&self.client, sandbox).await { + Ok(name) => name, + Err(e) => { + let _ = self.client.remove_volume(&vol_name).await; + return Err(e); + } + }; + let proxy_auth_secret_name = + match create_sandbox_proxy_auth_secret(&self.client, &self.config, sandbox) + .await + { + Ok(name) => name, + Err(e) => { + let _ = self.client.remove_volume(&vol_name).await; + if let Some(secret) = token_secret_name.as_deref() { + cleanup_sandbox_token_secret(&self.client, secret).await; + } + return Err(e); + } + }; + Ok((token_secret_name, proxy_auth_secret_name)) } - }; - let proxy_auth_secret_name = - match create_sandbox_proxy_auth_secret(&self.client, &self.config, sandbox).await { - Ok(name) => name, - Err(e) => { - let _ = self.client.remove_volume(&vol_name).await; - if let Some(secret) = token_secret_name.as_deref() { - cleanup_sandbox_token_secret(&self.client, secret).await; - } - return Err(e); - } - }; + .await; + phase_status.finish(result) + } + .instrument(tracing::info_span!( + "podman.prepare_storage", + otel.name = "podman.prepare_storage", + otel.status_code = tracing::field::Empty, + volume.name = %vol_name, + )) + .await?; // Clean up the volume and both per-sandbox secrets on any failure past // this point. @@ -800,41 +846,93 @@ impl PodmanComputeDriver { } }; - // 3. Create container. - let gpu_devices = match self.resolve_gpu_cdi_devices( - validated.gpu_requirements, - &validated.driver_config, - CdiGpuDefaultSelector::next_device_ids, - ) { - Ok(devices) => devices, - Err(e) => { - cleanup_created().await; - return Err(e); - } - }; - let supervisor_bin_path = if userns_needs_extraction(self.config.userns.as_deref()) { - match extract_supervisor_bin(&self.client, &self.config).await { - Ok(path) => Some(path), - Err(e) => { - cleanup_created().await; - return Err(e); - } - } - } else { - None - }; + // Prepare and create the container. + let tls_secret_names = async { + let phase_status = openshell_otel::ErrorStatusGuard::current(); + let result = async { + let gpu_devices = match self.resolve_gpu_cdi_devices( + validated.gpu_requirements, + &validated.driver_config, + CdiGpuDefaultSelector::next_device_ids, + ) { + Ok(devices) => devices, + Err(e) => { + cleanup_created().await; + return Err(e); + } + }; + let supervisor_bin_path = if userns_needs_extraction(self.config.userns.as_deref()) + { + match extract_supervisor_bin(&self.client, &self.config).await { + Ok(path) => Some(path), + Err(e) => { + cleanup_created().await; + return Err(e); + } + } + } else { + None + }; + + let tls_secret_names = if userns_remaps_uids(self.config.userns.as_deref()) + && self.config.tls_enabled() + { + let names = container::tls_secret_names(&sandbox.id); + if let Err(e) = create_tls_secrets(&self.client, &self.config, &names).await { + cleanup_created().await; + return Err(e); + } + Some(names) + } else { + None + }; - let tls_secret_names = - if userns_remaps_uids(self.config.userns.as_deref()) && self.config.tls_enabled() { - let names = container::tls_secret_names(&sandbox.id); - if let Err(e) = create_tls_secrets(&self.client, &self.config, &names).await { + let cleanup_all = || async { cleanup_created().await; - return Err(e); + if let Some(names) = &tls_secret_names { + cleanup_tls_secrets(&self.client, names).await; + } + }; + + let spec = match container::build_container_spec_for_image( + sandbox, + &self.config, + token_secret_name.as_deref(), + gpu_devices.as_deref(), + &image, + &immutable_image_id, + &image_user, + supervisor_bin_path.as_deref(), + tls_secret_names.as_ref(), + ) { + Ok(spec) => spec, + Err(e) => { + cleanup_all().await; + return Err(e); + } + }; + match self.client.create_container(&spec).await { + Ok(_) => Ok(tls_secret_names), + Err(PodmanApiError::Conflict(_)) => { + cleanup_all().await; + Err(ComputeDriverError::AlreadyExists) + } + Err(e) => { + cleanup_all().await; + Err(ComputeDriverError::from(e)) + } } - Some(names) - } else { - None - }; + } + .await; + phase_status.finish(result) + } + .instrument(tracing::info_span!( + "podman.prepare_container", + otel.name = "podman.prepare_container", + otel.status_code = tracing::field::Empty, + container.name = %name, + )) + .await?; let cleanup_all = || async { cleanup_created().await; @@ -843,37 +941,24 @@ impl PodmanComputeDriver { } }; - let spec = match container::build_container_spec_for_image( - sandbox, - &self.config, - token_secret_name.as_deref(), - gpu_devices.as_deref(), - image, - &inspected_image.id, - image_user, - supervisor_bin_path.as_deref(), - tls_secret_names.as_ref(), - ) { - Ok(spec) => spec, - Err(e) => { - cleanup_all().await; - return Err(e); - } - }; - match self.client.create_container(&spec).await { - Ok(_) => {} - Err(PodmanApiError::Conflict(_)) => { - cleanup_all().await; - return Err(ComputeDriverError::AlreadyExists); - } - Err(e) => { - cleanup_all().await; - return Err(ComputeDriverError::from(e)); - } + // Start container. + let start_result = async { + let phase_status = openshell_otel::ErrorStatusGuard::current(); + let result = self + .client + .start_container(&name) + .await + .map_err(ComputeDriverError::from); + phase_status.finish(result) } - - // 5. Start container. - if let Err(e) = self.client.start_container(&name).await { + .instrument(tracing::info_span!( + "podman.start_container", + otel.name = "podman.start_container", + otel.status_code = tracing::field::Empty, + container.name = %name, + )) + .await; + if let Err(e) = start_result { warn!( sandbox_name = %sandbox.name, error = %e, @@ -884,7 +969,7 @@ impl PodmanComputeDriver { .remove_container(&name, self.config.stop_timeout_secs) .await; cleanup_all().await; - return Err(ComputeDriverError::from(e)); + return Err(e); } info!( @@ -893,7 +978,7 @@ impl PodmanComputeDriver { "Sandbox container started" ); - Ok(()) + span_status.finish(Ok(())) } /// Find the Podman container ID for a sandbox by its sandbox ID using label lookup. @@ -948,44 +1033,69 @@ impl PodmanComputeDriver { } /// Stop a sandbox container without deleting it. + #[tracing::instrument( + name = "podman.stop_sandbox", + skip(self), + fields( + otel.name = "podman.stop_sandbox", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + ) + )] pub async fn stop_sandbox(&self, sandbox_id: &str) -> Result<(), ComputeDriverError> { + let span_status = openshell_otel::ErrorStatusGuard::current(); let container = self .find_container(sandbox_id) .await? .ok_or(ComputeDriverError::NotFound)?; let container_id = container.id; if container.state == "stopping" { - return self + let result = self .wait_for_container_stopped(sandbox_id, &container_id) .await; + return span_status.finish(result); } if container.state != "running" { - return Ok(()); + return span_status.finish(Ok(())); } info!(sandbox_id = %sandbox_id, container = %container_id, "Stopping sandbox container"); - self.client - .stop_container(&container_id, self.config.stop_timeout_secs) - .await - .map_err(ComputeDriverError::from)?; + let result = async { + self.client + .stop_container(&container_id, self.config.stop_timeout_secs) + .await + .map_err(ComputeDriverError::from)?; - // Podman can return from the stop request before inspect reports the - // container as exited. If start runs during that interval, the exit - // event from the previous run can arrive after the gateway has moved - // the same sandbox to Starting, causing it to regress to Error. Wait - // for the terminal container state before allowing a restart. - self.wait_for_container_stopped(sandbox_id, &container_id) - .await + // Podman can return from the stop request before inspect reports the + // container as exited. If start runs during that interval, the exit + // event from the previous run can arrive after the gateway has moved + // the same sandbox to Starting, causing it to regress to Error. Wait + // for the terminal container state before allowing a restart. + self.wait_for_container_stopped(sandbox_id, &container_id) + .await + } + .await; + span_status.finish(result) } /// Start a previously stopped sandbox container. + #[tracing::instrument( + name = "podman.start_sandbox", + skip(self), + fields( + otel.name = "podman.start_sandbox", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + ) + )] pub async fn start_sandbox(&self, sandbox_id: &str) -> Result<(), ComputeDriverError> { + let span_status = openshell_otel::ErrorStatusGuard::current(); let container = self .find_container(sandbox_id) .await? .ok_or(ComputeDriverError::NotFound)?; if container.state == "running" { - return Ok(()); + return span_status.finish(Ok(())); } let container_id = container.id; info!(sandbox_id = %sandbox_id, container = %container_id, "Starting sandbox container"); @@ -1002,14 +1112,26 @@ impl PodmanComputeDriver { .map_err(ComputeDriverError::from)?; self.lifecycle_event_fences .record_previous_exit(sandbox_id, previous.state.finished_at.as_deref()); - self.client + let result = self + .client .start_container(&container_id) .await - .map_err(ComputeDriverError::from) + .map_err(ComputeDriverError::from); + span_status.finish(result) } /// Delete a sandbox container and its workspace volume. + #[tracing::instrument( + name = "podman.delete_sandbox", + skip(self), + fields( + otel.name = "podman.delete_sandbox", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + ) + )] pub async fn delete_sandbox(&self, sandbox_id: &str) -> Result { + let span_status = openshell_otel::ErrorStatusGuard::current(); if sandbox_id.is_empty() { return Err(ComputeDriverError::Precondition( "sandbox id is required".into(), @@ -1031,7 +1153,7 @@ impl PodmanComputeDriver { .await; cleanup_tls_secrets(&self.client, &container::tls_secret_names(sandbox_id)).await; self.lifecycle_event_fences.remove(sandbox_id); - return Ok(false); + return span_status.finish(Ok(false)); }; info!(sandbox_id = %sandbox_id, container = %container_id, "Deleting sandbox container"); @@ -1067,7 +1189,7 @@ impl PodmanComputeDriver { cleanup_tls_secrets(&self.client, &container::tls_secret_names(sandbox_id)).await; self.lifecycle_event_fences.remove(sandbox_id); - Ok(container_existed) + span_status.finish(Ok(container_existed)) } /// Check whether a sandbox container exists. @@ -1611,6 +1733,215 @@ mod tests { let _ = fs::remove_file(socket); } + #[tokio::test] + async fn stop_sandbox_exports_a_podman_operation_span() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let (socket_path, _requests, handle) = spawn_podman_stub( + "trace-stop", + vec![ + StubResponse::new(StatusCode::OK, r#"[{"Id":"ctr-1","State":"running"}]"#), + StubResponse::new(StatusCode::NO_CONTENT, ""), + StubResponse::new( + StatusCode::OK, + r#"{"Id":"ctr-1","Name":"sandbox","State":{"Status":"exited","Running":false,"FinishedAt":"2026-08-12T16:39:13Z"},"Config":{}}"#, + ), + ], + ); + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)); + + test_driver(socket_path.clone()) + .stop_sandbox("sandbox-1") + .with_subscriber(subscriber) + .await + .expect("stop should succeed"); + handle.await.expect("stub should finish"); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "podman.stop_sandbox") + .expect("stop operation should be exported"); + assert_eq!( + span.attributes + .iter() + .find(|attribute| attribute.key.as_str() == "sandbox.id") + .map(|attribute| attribute.value.to_string()) + .as_deref(), + Some("sandbox-1") + ); + provider.shutdown().unwrap(); + let _ = fs::remove_file(socket_path); + } + + #[tokio::test] + async fn create_sandbox_exports_nested_preparation_spans() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let (socket_path, _requests, handle) = spawn_podman_stub( + "trace-create", + vec![ + StubResponse::new(StatusCode::OK, "{}"), + StubResponse::new(StatusCode::OK, "{}"), + StubResponse::new( + StatusCode::OK, + r#"{"Id":"sha256:sandbox","Config":{"User":"1234:1235"}}"#, + ), + StubResponse::new(StatusCode::CREATED, "{}"), + StubResponse::new(StatusCode::CREATED, "{}"), + StubResponse::new(StatusCode::NO_CONTENT, ""), + ], + ); + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)); + + test_driver(socket_path.clone()) + .create_sandbox(&plain_sandbox("sandbox-trace", "demo")) + .with_subscriber(subscriber) + .await + .expect("create should succeed"); + handle.await.expect("stub should finish"); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let create = spans + .iter() + .find(|span| span.name == "podman.create_sandbox") + .expect("create operation should be exported"); + for name in [ + "podman.prepare_images", + "podman.prepare_storage", + "podman.prepare_container", + "podman.start_container", + ] { + let child = spans + .iter() + .find(|span| span.name == name) + .unwrap_or_else(|| panic!("{name} should be exported")); + assert_eq!( + child.parent_span_id, + create.span_context.span_id(), + "{name}" + ); + } + provider.shutdown().unwrap(); + let _ = fs::remove_file(socket_path); + } + + #[tokio::test] + async fn prepare_images_span_covers_and_marks_sandbox_image_pull_failure() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let (socket_path, _requests, handle) = spawn_podman_stub( + "trace-image-failure", + vec![ + StubResponse::new(StatusCode::OK, "{}"), + StubResponse::new(StatusCode::INTERNAL_SERVER_ERROR, "pull failed"), + ], + ); + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)); + + test_driver(socket_path.clone()) + .create_sandbox(&plain_sandbox("sandbox-trace", "demo")) + .with_subscriber(subscriber) + .await + .expect_err("sandbox image pull should fail"); + handle.await.expect("stub should finish"); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let phase = spans + .iter() + .find(|span| span.name == "podman.prepare_images") + .expect("image preparation should be exported"); + assert!(matches!( + phase.status, + opentelemetry::trace::Status::Error { .. } + )); + provider.shutdown().unwrap(); + let _ = fs::remove_file(socket_path); + } + + #[tokio::test] + async fn start_and_delete_export_podman_operation_spans() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)); + + let (start_socket, _requests, start_handle) = spawn_podman_stub( + "trace-start", + vec![ + StubResponse::new(StatusCode::OK, r#"[{"Id":"ctr-1","State":"stopped"}]"#), + StubResponse::new( + StatusCode::OK, + r#"{"Id":"ctr-1","Name":"sandbox","State":{"Status":"exited","Running":false,"FinishedAt":"2026-08-12T16:39:13Z"},"Config":{}}"#, + ), + StubResponse::new(StatusCode::NO_CONTENT, ""), + ], + ); + test_driver(start_socket.clone()) + .start_sandbox("sandbox-1") + .with_subscriber(subscriber) + .await + .expect("start should succeed"); + start_handle.await.expect("start stub should finish"); + + let (delete_socket, _requests, delete_handle) = spawn_podman_stub( + "trace-delete", + vec![ + StubResponse::new(StatusCode::OK, "[]"), + StubResponse::new(StatusCode::NO_CONTENT, ""), + ], + ); + let subscriber = tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)); + test_driver(delete_socket.clone()) + .delete_sandbox("sandbox-1") + .with_subscriber(subscriber) + .await + .expect("delete should succeed"); + delete_handle.await.expect("delete stub should finish"); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + assert!(spans.iter().any(|span| span.name == "podman.start_sandbox")); + assert!( + spans + .iter() + .any(|span| span.name == "podman.delete_sandbox") + ); + provider.shutdown().unwrap(); + let _ = fs::remove_file(start_socket); + let _ = fs::remove_file(delete_socket); + } + #[test] fn validate_gpu_request_accepts_gpu_count_request_shape() { let gpu = GpuResourceRequirements { count: Some(2) }; diff --git a/crates/openshell-driver-podman/src/grpc.rs b/crates/openshell-driver-podman/src/grpc.rs index 19d0b55254..7a7275a8a6 100644 --- a/crates/openshell-driver-podman/src/grpc.rs +++ b/crates/openshell-driver-podman/src/grpc.rs @@ -14,20 +14,114 @@ use openshell_core::proto::compute::v1::{ ValidateSandboxCreateRequest, ValidateSandboxCreateResponse, WatchSandboxesEvent, WatchSandboxesRequest, compute_driver_server::ComputeDriver, }; +use std::future::Future; use std::pin::Pin; +use std::task::{Context, Poll}; use tonic::{Request, Response, Status}; +use tracing::Instrument as _; use crate::PodmanComputeDriver; +type ComputeDriverWatchStream = + Pin> + Send + 'static>>; + +struct TracedWatchStream { + inner: ComputeDriverWatchStream, + span: tracing::Span, +} + +impl TracedWatchStream { + fn new(inner: ComputeDriverWatchStream, span: tracing::Span) -> Self { + Self { inner, span } + } +} + +impl Stream for TracedWatchStream { + type Item = Result; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + let span = self.span.clone(); + let _entered = span.enter(); + let result = self.inner.as_mut().poll_next(cx); + match &result { + Poll::Ready(Some(Err(status))) => { + openshell_otel::mark_error(&self.span); + self.span + .record("rpc.grpc.status_code", status.code() as i32); + } + Poll::Ready(None) => { + self.span + .record("rpc.grpc.status_code", tonic::Code::Ok as i32); + } + Poll::Pending | Poll::Ready(Some(Ok(_))) => {} + } + result + } +} + #[derive(Debug, Clone)] pub struct ComputeDriverService { driver: PodmanComputeDriver, + trace_in_process_rpc: bool, } impl ComputeDriverService { #[must_use] pub fn new(driver: PodmanComputeDriver) -> Self { - Self { driver } + Self { + driver, + trace_in_process_rpc: false, + } + } + + #[must_use] + pub fn new_in_process(driver: PodmanComputeDriver) -> Self { + Self { + driver, + trace_in_process_rpc: true, + } + } + + fn in_process_rpc_span( + &self, + operation: &'static str, + method: &'static str, + ) -> Option { + self.trace_in_process_rpc.then(|| { + tracing::info_span!( + target: "openshell_driver_podman::otel_tracing", + "driver_rpc", + otel.name = operation, + otel.kind = "server", + otel.status_code = tracing::field::Empty, + rpc.system = "grpc", + rpc.service = "openshell.compute.v1.ComputeDriver", + rpc.method = method, + rpc.grpc.status_code = tracing::field::Empty, + ) + }) + } + + async fn trace_rpc( + &self, + operation: &'static str, + method: &'static str, + future: impl Future>, + ) -> Result { + let Some(span) = self.in_process_rpc_span(operation, method) else { + return future.await; + }; + let result = future.instrument(span.clone()).await; + match &result { + Ok(_) => { + span.record("rpc.grpc.status_code", tonic::Code::Ok as i32); + } + Err(status) => { + openshell_otel::mark_error(&span); + span.record("rpc.grpc.status_code", status.code() as i32); + } + } + result } } @@ -37,153 +131,207 @@ impl ComputeDriver for ComputeDriverService { &self, _request: Request, ) -> Result, Status> { - self.driver - .capabilities() - .map(Response::new) - .map_err(Status::from) + self.trace_rpc("driver.get_capabilities", "get_capabilities", async { + self.driver + .capabilities() + .map(Response::new) + .map_err(Status::from) + }) + .await } async fn get_gateway_listener_requirements( &self, _request: Request, ) -> Result, Status> { - Ok(Response::new(GetGatewayListenerRequirementsResponse { - requirements: self - .driver - .gateway_listener_requirements() - .map_err(Status::from)?, - })) + self.trace_rpc( + "driver.get_gateway_listener_requirements", + "get_gateway_listener_requirements", + async { + Ok(Response::new(GetGatewayListenerRequirementsResponse { + requirements: self + .driver + .gateway_listener_requirements() + .map_err(Status::from)?, + })) + }, + ) + .await } async fn validate_sandbox_create( &self, request: Request, ) -> Result, Status> { - let sandbox = request - .into_inner() - .sandbox - .ok_or_else(|| Status::invalid_argument("sandbox is required"))?; - self.driver - .validate_sandbox_create(&sandbox) - .await - .map_err(Status::from)?; - Ok(Response::new(ValidateSandboxCreateResponse {})) + self.trace_rpc( + "driver.validate_sandbox_create", + "validate_sandbox_create", + async { + let sandbox = request + .into_inner() + .sandbox + .ok_or_else(|| Status::invalid_argument("sandbox is required"))?; + self.driver + .validate_sandbox_create(&sandbox) + .await + .map_err(Status::from)?; + Ok(Response::new(ValidateSandboxCreateResponse {})) + }, + ) + .await } async fn get_sandbox( &self, request: Request, ) -> Result, Status> { - let request = request.into_inner(); - if request.sandbox_id.is_empty() { - return Err(Status::invalid_argument("sandbox_id is required")); - } - - let sandbox = self - .driver - .get_sandbox(&request.sandbox_id) - .await - .map_err(Status::from)? - .ok_or_else(|| Status::not_found("sandbox not found"))?; - - Ok(Response::new(GetSandboxResponse { - sandbox: Some(sandbox), - })) + self.trace_rpc("driver.get_sandbox", "get_sandbox", async { + let request = request.into_inner(); + if request.sandbox_id.is_empty() { + return Err(Status::invalid_argument("sandbox_id is required")); + } + let sandbox = self + .driver + .get_sandbox(&request.sandbox_id) + .await + .map_err(Status::from)? + .ok_or_else(|| Status::not_found("sandbox not found"))?; + Ok(Response::new(GetSandboxResponse { + sandbox: Some(sandbox), + })) + }) + .await } async fn list_sandboxes( &self, _request: Request, ) -> Result, Status> { - let sandboxes = self.driver.list_sandboxes().await.map_err(Status::from)?; - Ok(Response::new(ListSandboxesResponse { sandboxes })) + self.trace_rpc("driver.list_sandboxes", "list_sandboxes", async { + let sandboxes = self.driver.list_sandboxes().await.map_err(Status::from)?; + Ok(Response::new(ListSandboxesResponse { sandboxes })) + }) + .await } async fn create_sandbox( &self, request: Request, ) -> Result, Status> { - let sandbox = request - .into_inner() - .sandbox - .ok_or_else(|| Status::invalid_argument("sandbox is required"))?; - self.driver - .create_sandbox(&sandbox) - .await - .map_err(Status::from)?; - Ok(Response::new(CreateSandboxResponse {})) + self.trace_rpc("driver.create_sandbox", "create_sandbox", async { + let sandbox = request + .into_inner() + .sandbox + .ok_or_else(|| Status::invalid_argument("sandbox is required"))?; + self.driver + .create_sandbox(&sandbox) + .await + .map_err(Status::from)?; + Ok(Response::new(CreateSandboxResponse {})) + }) + .await } async fn stop_sandbox( &self, request: Request, ) -> Result, Status> { - let request = request.into_inner(); - if request.sandbox_id.is_empty() { - return Err(Status::invalid_argument("sandbox_id is required")); - } - self.driver - .stop_sandbox(&request.sandbox_id) - .await - .map_err(Status::from)?; - Ok(Response::new(StopSandboxResponse {})) + self.trace_rpc("driver.stop_sandbox", "stop_sandbox", async { + let request = request.into_inner(); + if request.sandbox_id.is_empty() { + return Err(Status::invalid_argument("sandbox_id is required")); + } + self.driver + .stop_sandbox(&request.sandbox_id) + .await + .map_err(Status::from)?; + Ok(Response::new(StopSandboxResponse {})) + }) + .await } async fn start_sandbox( &self, request: Request, ) -> Result, Status> { - let request = request.into_inner(); - if request.sandbox_id.is_empty() { - return Err(Status::invalid_argument("sandbox_id is required")); - } - self.driver - .start_sandbox(&request.sandbox_id) - .await - .map_err(Status::from)?; - Ok(Response::new(StartSandboxResponse {})) + self.trace_rpc("driver.start_sandbox", "start_sandbox", async { + let request = request.into_inner(); + if request.sandbox_id.is_empty() { + return Err(Status::invalid_argument("sandbox_id is required")); + } + self.driver + .start_sandbox(&request.sandbox_id) + .await + .map_err(Status::from)?; + Ok(Response::new(StartSandboxResponse {})) + }) + .await } async fn delete_sandbox( &self, request: Request, ) -> Result, Status> { - let request = request.into_inner(); - if request.sandbox_id.is_empty() { - return Err(Status::invalid_argument("sandbox_id is required")); - } - let deleted = self - .driver - .delete_sandbox(&request.sandbox_id) - .await - .map_err(Status::from)?; - Ok(Response::new(DeleteSandboxResponse { deleted })) + self.trace_rpc("driver.delete_sandbox", "delete_sandbox", async { + let request = request.into_inner(); + if request.sandbox_id.is_empty() { + return Err(Status::invalid_argument("sandbox_id is required")); + } + let deleted = self + .driver + .delete_sandbox(&request.sandbox_id) + .await + .map_err(Status::from)?; + Ok(Response::new(DeleteSandboxResponse { deleted })) + }) + .await } - type WatchSandboxesStream = - Pin> + Send + 'static>>; + type WatchSandboxesStream = ComputeDriverWatchStream; async fn watch_sandboxes( &self, _request: Request, ) -> Result, Status> { - let stream = self.driver.watch_sandboxes().await.map_err(Status::from)?; - let stream = stream.map(|item| item.map_err(|err| Status::internal(err.to_string()))); - Ok(Response::new(Box::pin(stream))) + let create_stream = async { + let stream = self.driver.watch_sandboxes().await.map_err(Status::from)?; + let stream = stream.map(|item| item.map_err(|err| Status::internal(err.to_string()))); + Ok::(Box::pin(stream)) + }; + let Some(span) = self.in_process_rpc_span("driver.watch_sandboxes", "watch_sandboxes") + else { + return create_stream.await.map(Response::new); + }; + match create_stream.instrument(span.clone()).await { + Ok(stream) => Ok(Response::new(Box::pin(TracedWatchStream::new( + stream, span, + )))), + Err(status) => { + openshell_otel::mark_error(&span); + span.record("rpc.grpc.status_code", status.code() as i32); + Err(status) + } + } } async fn ensure_workspace( &self, _request: Request, ) -> Result, Status> { - Ok(Response::new(EnsureWorkspaceResponse {})) + self.trace_rpc("driver.ensure_workspace", "ensure_workspace", async { + Ok(Response::new(EnsureWorkspaceResponse {})) + }) + .await } async fn delete_workspace( &self, _request: Request, ) -> Result, Status> { - Ok(Response::new(DeleteWorkspaceResponse {})) + self.trace_rpc("driver.delete_workspace", "delete_workspace", async { + Ok(Response::new(DeleteWorkspaceResponse {})) + }) + .await } } @@ -197,6 +345,53 @@ mod tests { use openshell_core::ComputeDriverError; use std::path::PathBuf; + type TestDriverClient = + openshell_core::proto::compute::v1::compute_driver_client::ComputeDriverClient< + tonic::transport::Channel, + >; + + fn request_with_traceparent(message: T) -> Request { + let mut request = Request::new(message); + request.metadata_mut().insert( + "traceparent", + "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01" + .parse() + .unwrap(), + ); + request + } + + async fn standalone_traced_client() -> ( + TestDriverClient, + tokio::sync::oneshot::Sender<()>, + tokio::task::JoinHandle>, + ) { + use openshell_core::proto::compute::v1::compute_driver_server::ComputeDriverServer; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let (shutdown, shutdown_rx) = tokio::sync::oneshot::channel(); + let service = ComputeDriverService::new(PodmanComputeDriver::for_tests( + PodmanComputeConfig::default(), + )); + let server = tokio::spawn(async move { + tonic::transport::Server::builder() + .layer(crate::otel_tracing::compute_driver_rpc_layer()) + .add_service(ComputeDriverServer::new(service)) + .serve_with_incoming_shutdown( + tokio_stream::wrappers::TcpListenerStream::new(listener), + async { + let _ = shutdown_rx.await; + }, + ) + .await + }); + let client = TestDriverClient::connect(format!("http://{address}")) + .await + .unwrap(); + (client, shutdown, server) + } + #[test] fn precondition_driver_errors_map_to_failed_precondition_status() { let status: Status = @@ -218,6 +413,242 @@ mod tests { assert_eq!(status.code(), tonic::Code::NotFound); } + #[tokio::test] + async fn in_process_service_preserves_the_driver_rpc_server_boundary() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::{Instrument as _, instrument::WithSubscriber as _}; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let gateway_exporter = InMemorySpanExporterBuilder::new().build(); + let gateway_provider = SdkTracerProvider::builder() + .with_simple_exporter(gateway_exporter.clone()) + .build(); + let driver_exporter = InMemorySpanExporterBuilder::new().build(); + let driver_provider = SdkTracerProvider::builder() + .with_simple_exporter(driver_exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry() + .with(openshell_otel::layer_excluding_target_prefix( + &gateway_provider, + "gateway-test", + crate::otel_tracing::IN_PROCESS_TARGET_PREFIX, + )) + .with(crate::otel_tracing::in_process_layer(&driver_provider)); + let service = ComputeDriverService::new_in_process(PodmanComputeDriver::for_tests( + PodmanComputeConfig::default(), + )); + + async { + let gateway_span = tracing::info_span!(target: "openshell_server::compute", "driver", otel.name = "driver.get_capabilities", otel.kind = "client"); + ComputeDriver::get_capabilities(&service, Request::new(GetCapabilitiesRequest {})) + .instrument(gateway_span) + .await + } + .with_subscriber(subscriber) + .await + .expect("capabilities should succeed"); + gateway_provider.force_flush().unwrap(); + driver_provider.force_flush().unwrap(); + + let gateway_spans = gateway_exporter.get_finished_spans().unwrap(); + let driver_spans = driver_exporter.get_finished_spans().unwrap(); + let client = gateway_spans + .iter() + .find(|span| span.name == "driver.get_capabilities") + .unwrap(); + let server = driver_spans + .iter() + .find(|span| span.name == "driver.get_capabilities") + .expect("in-process server span"); + assert_eq!( + server.span_context.trace_id(), + client.span_context.trace_id() + ); + assert_eq!(server.parent_span_id, client.span_context.span_id()); + assert_eq!(server.span_kind, opentelemetry::trace::SpanKind::Server); + assert!(server.attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" + && attribute.value.to_string() == (tonic::Code::Ok as i32).to_string() + })); + gateway_provider.shutdown().unwrap(); + driver_provider.shutdown().unwrap(); + } + + #[tokio::test] + async fn standalone_rpc_layer_propagates_context_and_records_errors() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let dispatch = tracing::Dispatch::new( + tracing_subscriber::registry().with(crate::otel_tracing::layer(&provider)), + ); + let _dispatch = tracing::dispatcher::set_default(&dispatch); + let (mut client, shutdown, server) = standalone_traced_client().await; + + client + .get_capabilities(request_with_traceparent(GetCapabilitiesRequest {})) + .await + .expect("capabilities should succeed"); + client + .validate_sandbox_create(request_with_traceparent(ValidateSandboxCreateRequest { + sandbox: None, + })) + .await + .expect_err("missing sandbox should fail"); + drop(client); + shutdown.send(()).unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), server) + .await + .expect("standalone test server should stop") + .expect("standalone test server should not panic") + .expect("standalone test server should stop cleanly"); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let capabilities = spans + .iter() + .filter(|span| span.name == "driver.get_capabilities") + .collect::>(); + assert_eq!(capabilities.len(), 1, "exactly one RPC span is expected"); + assert_eq!( + capabilities[0].span_context.trace_id().to_string(), + "4bf92f3577b34da6a3ce929d0e0e4736" + ); + assert_eq!( + capabilities[0].parent_span_id.to_string(), + "00f067aa0ba902b7" + ); + assert_eq!( + capabilities[0].span_kind, + opentelemetry::trace::SpanKind::Server + ); + assert!(capabilities[0].attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" + && attribute.value.to_string() == (tonic::Code::Ok as i32).to_string() + })); + + let failed = spans + .iter() + .find(|span| span.name == "driver.validate_sandbox_create") + .expect("failed RPC span should be exported"); + assert!(matches!( + failed.status, + opentelemetry::trace::Status::Error { .. } + )); + assert!(failed.attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" + && attribute.value.to_string() == (tonic::Code::InvalidArgument as i32).to_string() + })); + provider.shutdown().unwrap(); + } + + #[tokio::test] + async fn in_process_stream_span_lives_until_stream_failure() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = + tracing_subscriber::registry().with(crate::otel_tracing::in_process_layer(&provider)); + + async { + let span = tracing::info_span!( + target: "openshell_driver_podman::otel_tracing", + "driver_rpc", + otel.name = "driver.watch_sandboxes", + otel.kind = "server", + otel.status_code = tracing::field::Empty, + rpc.grpc.status_code = tracing::field::Empty, + ); + let inner: ComputeDriverWatchStream = Box::pin(futures::stream::iter([Err( + Status::internal("watch failed"), + )])); + let mut stream = TracedWatchStream::new(inner, span); + + provider.force_flush().unwrap(); + assert!( + exporter.get_finished_spans().unwrap().is_empty(), + "server span must remain open while the response stream is alive" + ); + stream + .next() + .await + .expect("stream item") + .expect_err("stream should fail"); + drop(stream); + } + .with_subscriber(subscriber) + .await; + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "driver.watch_sandboxes") + .expect("watch server span should be exported when the stream ends"); + assert!(matches!( + span.status, + opentelemetry::trace::Status::Error { .. } + )); + provider.shutdown().unwrap(); + } + + #[tokio::test] + async fn in_process_stream_records_ok_when_stream_completes() { + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing::instrument::WithSubscriber as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = crate::otel_tracing::test_lock().await; + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = + tracing_subscriber::registry().with(crate::otel_tracing::in_process_layer(&provider)); + + async { + let span = tracing::info_span!( + target: "openshell_driver_podman::otel_tracing", + "driver_rpc", + otel.name = "driver.watch_sandboxes", + otel.kind = "server", + otel.status_code = tracing::field::Empty, + rpc.grpc.status_code = tracing::field::Empty, + ); + let inner: ComputeDriverWatchStream = Box::pin(futures::stream::empty()); + let mut stream = TracedWatchStream::new(inner, span); + + assert!(stream.next().await.is_none()); + drop(stream); + } + .with_subscriber(subscriber) + .await; + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "driver.watch_sandboxes") + .expect("watch server span should be exported when the stream completes"); + assert!(span.attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" + && attribute.value.to_string() == (tonic::Code::Ok as i32).to_string() + })); + provider.shutdown().unwrap(); + } + fn test_service(socket_path: PathBuf) -> ComputeDriverService { let config = PodmanComputeConfig { socket_path: Some(socket_path), diff --git a/crates/openshell-driver-podman/src/lib.rs b/crates/openshell-driver-podman/src/lib.rs index 5847a10ea6..7a78cf5d92 100644 --- a/crates/openshell-driver-podman/src/lib.rs +++ b/crates/openshell-driver-podman/src/lib.rs @@ -6,6 +6,7 @@ pub mod config; pub(crate) mod container; pub mod driver; pub mod grpc; +pub mod otel_tracing; #[cfg(test)] pub(crate) mod test_utils; pub(crate) mod watcher; diff --git a/crates/openshell-driver-podman/src/main.rs b/crates/openshell-driver-podman/src/main.rs index 405deb93a6..81d4254d10 100644 --- a/crates/openshell-driver-podman/src/main.rs +++ b/crates/openshell-driver-podman/src/main.rs @@ -3,10 +3,12 @@ use clap::Parser; use miette::{IntoDiagnostic, Result}; +use std::future::Future; use std::net::SocketAddr; use std::path::PathBuf; use tracing::info; use tracing_subscriber::EnvFilter; +use tracing_subscriber::prelude::*; use openshell_core::VERSION; use openshell_core::proto::compute::v1::compute_driver_server::ComputeDriverServer; @@ -14,6 +16,7 @@ use openshell_driver_podman::config::{ DEFAULT_NETWORK_NAME, DEFAULT_PODMAN_STOP_TIMEOUT_SECS, DEFAULT_SANDBOX_PIDS_LIMIT, ImagePullPolicy, }; +use openshell_driver_podman::otel_tracing::compute_driver_rpc_layer; use openshell_driver_podman::{ComputeDriverService, PodmanComputeConfig, PodmanComputeDriver}; #[derive(Parser)] @@ -30,6 +33,9 @@ struct Args { #[arg(long, env = "OPENSHELL_LOG_LEVEL", default_value = "info")] log_level: String, + #[arg(long, env = "OPENSHELL_OTLP_ENDPOINT")] + otlp_endpoint: Option, + /// Path to the Podman API Unix socket. #[arg(long, env = "OPENSHELL_PODMAN_SOCKET")] podman_socket: Option, @@ -153,11 +159,22 @@ struct Args { #[tokio::main] async fn main() -> Result<()> { let args = Args::parse(); - tracing_subscriber::fmt() - .with_env_filter( - EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(&args.log_level)), + let (tracer_provider, setup_error) = + openshell_driver_podman::otel_tracing::provider_for(args.otlp_endpoint.as_deref()); + tracing_subscriber::registry() + .with(EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new(&args.log_level))) + .with(tracing_subscriber::fmt::layer()) + .with( + tracer_provider + .as_ref() + .map(openshell_driver_podman::otel_tracing::layer), ) .init(); + if let Some(error) = setup_error { + tracing::error!(%error, "OTLP exporting could not be started"); + } else if let Some(endpoint) = &args.otlp_endpoint { + info!(endpoint, "OTLP exporting enabled"); + } let driver = PodmanComputeDriver::new(PodmanComputeConfig { socket_path: args.podman_socket, @@ -192,12 +209,80 @@ async fn main() -> Result<()> { .into_diagnostic()?; info!(address = %args.bind_address, "Starting Podman compute driver"); - tonic::transport::Server::builder() + let result = tonic::transport::Server::builder() + .layer(compute_driver_rpc_layer()) .add_service(ComputeDriverServer::new(ComputeDriverService::new(driver))) .serve_with_shutdown(args.bind_address, async { - tokio::signal::ctrl_c().await.ok(); + shutdown_signal().await; info!("Received shutdown signal, draining in-flight requests"); }) .await - .into_diagnostic() + .into_diagnostic(); + if let Some(provider) = &tracer_provider + && let Err(error) = provider.shutdown() + { + tracing::warn!(%error, "OTLP tracer provider shutdown failed"); + } + result +} + +async fn select_shutdown_signal( + ctrl_c: impl Future, + terminate: impl Future, +) { + tokio::select! { + () = ctrl_c => {} + () = terminate => {} + } +} + +async fn ctrl_c_signal() { + if let Err(error) = tokio::signal::ctrl_c().await { + tracing::warn!(%error, "Failed to install Ctrl-C signal handler"); + std::future::pending::<()>().await; + } +} + +#[cfg(unix)] +async fn terminate_signal() { + let Ok(mut signal) = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate()) + else { + tracing::warn!("Failed to install SIGTERM signal handler"); + std::future::pending::<()>().await; + return; + }; + let _ = signal.recv().await; +} + +async fn shutdown_signal() { + #[cfg(unix)] + select_shutdown_signal(ctrl_c_signal(), terminate_signal()).await; + + #[cfg(not(unix))] + ctrl_c_signal().await; +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn shutdown_completes_when_termination_signal_arrives() { + select_shutdown_signal(std::future::pending(), std::future::ready(())).await; + } + + #[test] + fn accepts_gateway_otlp_endpoint() { + let args = Args::try_parse_from([ + "openshell-driver-podman", + "--otlp-endpoint", + "http://collector.internal:4317", + ]) + .expect("OTLP endpoint should be accepted"); + + assert_eq!( + args.otlp_endpoint.as_deref(), + Some("http://collector.internal:4317") + ); + } } diff --git a/crates/openshell-driver-podman/src/otel_tracing.rs b/crates/openshell-driver-podman/src/otel_tracing.rs new file mode 100644 index 0000000000..9d8f494276 --- /dev/null +++ b/crates/openshell-driver-podman/src/otel_tracing.rs @@ -0,0 +1,299 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! OpenTelemetry trace exporting for the Podman compute driver. + +use http::Request; +use openshell_otel::{ + HeaderMapExtractor, OtlpTraceConfig, RecordGrpcFailure, RecordGrpcStatus, SdkTracerProvider, + ServiceName, SetupError, +}; +use opentelemetry::propagation::TextMapPropagator as _; +use opentelemetry::trace::TraceContextExt as _; +use opentelemetry_sdk::propagation::TraceContextPropagator; +use tower_http::trace::{GrpcMakeClassifier, MakeSpan, TraceLayer}; +use tracing::{Span, Subscriber}; +use tracing_opentelemetry::OpenTelemetrySpanExt as _; +use tracing_subscriber::registry::LookupSpan; + +const SERVICE_NAME: &str = "openshell-driver-podman"; +const INSTRUMENTATION_SCOPE: &str = "openshell-driver-podman"; +const COMPUTE_DRIVER_SERVICE: &str = "openshell.compute.v1.ComputeDriver"; +pub const IN_PROCESS_TARGET_PREFIX: &str = "openshell_driver_podman"; + +pub fn compute_driver_rpc_layer() -> TraceLayer< + GrpcMakeClassifier, + ComputeDriverRpcSpan, + (), + RecordGrpcStatus, + (), + RecordGrpcStatus, + RecordGrpcFailure, +> { + TraceLayer::new_for_grpc() + .make_span_with(ComputeDriverRpcSpan) + .on_request(()) + .on_response(RecordGrpcStatus) + .on_body_chunk(()) + .on_eos(RecordGrpcStatus) + .on_failure(RecordGrpcFailure) +} + +#[derive(Debug, Clone, Copy)] +pub struct ComputeDriverRpcSpan; + +impl MakeSpan for ComputeDriverRpcSpan { + fn make_span(&mut self, request: &Request) -> Span { + let (operation, method) = compute_driver_rpc_operation(request.uri().path()); + let span = tracing::info_span!( + "driver_rpc", + otel.name = operation, + otel.kind = "server", + otel.status_code = tracing::field::Empty, + rpc.system = "grpc", + rpc.service = COMPUTE_DRIVER_SERVICE, + rpc.method = method, + rpc.grpc.status_code = tracing::field::Empty, + ); + let parent = TraceContextPropagator::new().extract_with_context( + &opentelemetry::Context::new(), + &HeaderMapExtractor::new(request.headers()), + ); + if parent.span().span_context().is_valid() { + let _ = span.set_parent(parent); + } + span + } +} + +pub(crate) fn compute_driver_rpc_operation(path: &str) -> (&'static str, &'static str) { + match path.rsplit('/').next() { + Some("GetCapabilities") => ("driver.get_capabilities", "get_capabilities"), + Some("GetGatewayListenerRequirements") => ( + "driver.get_gateway_listener_requirements", + "get_gateway_listener_requirements", + ), + Some("ValidateSandboxCreate") => { + ("driver.validate_sandbox_create", "validate_sandbox_create") + } + Some("CreateSandbox") => ("driver.create_sandbox", "create_sandbox"), + Some("GetSandbox") => ("driver.get_sandbox", "get_sandbox"), + Some("ListSandboxes") => ("driver.list_sandboxes", "list_sandboxes"), + Some("StopSandbox") => ("driver.stop_sandbox", "stop_sandbox"), + Some("StartSandbox") => ("driver.start_sandbox", "start_sandbox"), + Some("DeleteSandbox") => ("driver.delete_sandbox", "delete_sandbox"), + Some("WatchSandboxes") => ("driver.watch_sandboxes", "watch_sandboxes"), + Some("EnsureWorkspace") => ("driver.ensure_workspace", "ensure_workspace"), + Some("DeleteWorkspace") => ("driver.delete_workspace", "delete_workspace"), + _ => ("driver.unknown", "unknown"), + } +} + +#[must_use] +pub fn provider_for(endpoint: Option<&str>) -> (Option, Option) { + openshell_otel::provider_for(endpoint.map(|endpoint| OtlpTraceConfig { + endpoint, + service_name: ServiceName::Fixed(SERVICE_NAME), + service_version: Some(openshell_core::VERSION), + resource_attributes: Vec::new(), + })) +} + +pub fn layer(provider: &SdkTracerProvider) -> openshell_otel::OtlpLayer +where + S: Subscriber + for<'span> LookupSpan<'span>, +{ + openshell_otel::layer(provider, INSTRUMENTATION_SCOPE) +} + +pub fn in_process_layer(provider: &SdkTracerProvider) -> openshell_otel::TargetOtlpLayer +where + S: Subscriber + for<'span> LookupSpan<'span>, +{ + openshell_otel::layer_for_target_prefix( + provider, + INSTRUMENTATION_SCOPE, + IN_PROCESS_TARGET_PREFIX, + ) +} + +#[cfg(test)] +pub(crate) async fn test_lock() -> tokio::sync::MutexGuard<'static, ()> { + static LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); + static INITIALIZED: std::sync::LazyLock<()> = std::sync::LazyLock::new(|| { + tracing::subscriber::set_global_default(tracing_subscriber::registry()) + .expect("test tracing subscriber installs once"); + }); + + let guard = LOCK.lock().await; + std::sync::LazyLock::force(&INITIALIZED); + guard +} + +#[cfg(test)] +mod tests { + use std::sync::{Arc, Mutex}; + + use opentelemetry_proto::tonic::collector::trace::v1::{ + ExportTraceServiceRequest, ExportTraceServiceResponse, + trace_service_server::{TraceService, TraceServiceServer}, + }; + use opentelemetry_proto::tonic::trace::v1::Span; + use tracing_subscriber::layer::SubscriberExt as _; + + #[derive(Default)] + struct Received { + spans: Vec, + service_names: Vec, + } + + #[derive(Clone)] + struct Collector { + received: Arc>, + exported: Arc, + } + + #[tonic::async_trait] + impl TraceService for Collector { + async fn export( + &self, + request: tonic::Request, + ) -> Result, tonic::Status> { + let mut received = self.received.lock().unwrap(); + for resource_span in request.into_inner().resource_spans { + if let Some(resource) = resource_span.resource { + received.service_names.extend( + resource + .attributes + .into_iter() + .filter(|attribute| attribute.key == "service.name") + .filter_map(|attribute| attribute.value) + .filter_map(|value| value.value) + .filter_map(|value| match value { + opentelemetry_proto::tonic::common::v1::any_value::Value::StringValue(value) => Some(value), + _ => None, + }), + ); + } + for scope_span in resource_span.scope_spans { + received.spans.extend(scope_span.spans); + } + } + drop(received); + self.exported.notify_one(); + Ok(tonic::Response::new(ExportTraceServiceResponse::default())) + } + } + + #[test] + fn compute_driver_rpc_names_are_explicitly_mapped_and_schema_bounded() { + for (rpc, operation, method) in [ + ( + "GetCapabilities", + "driver.get_capabilities", + "get_capabilities", + ), + ( + "GetGatewayListenerRequirements", + "driver.get_gateway_listener_requirements", + "get_gateway_listener_requirements", + ), + ( + "ValidateSandboxCreate", + "driver.validate_sandbox_create", + "validate_sandbox_create", + ), + ("CreateSandbox", "driver.create_sandbox", "create_sandbox"), + ("GetSandbox", "driver.get_sandbox", "get_sandbox"), + ("ListSandboxes", "driver.list_sandboxes", "list_sandboxes"), + ("StopSandbox", "driver.stop_sandbox", "stop_sandbox"), + ("StartSandbox", "driver.start_sandbox", "start_sandbox"), + ("DeleteSandbox", "driver.delete_sandbox", "delete_sandbox"), + ( + "WatchSandboxes", + "driver.watch_sandboxes", + "watch_sandboxes", + ), + ( + "EnsureWorkspace", + "driver.ensure_workspace", + "ensure_workspace", + ), + ( + "DeleteWorkspace", + "driver.delete_workspace", + "delete_workspace", + ), + ] { + assert_eq!( + super::compute_driver_rpc_operation(&format!( + "/openshell.compute.v1.ComputeDriver/{rpc}" + )), + (operation, method), + ); + } + assert_eq!( + super::compute_driver_rpc_operation( + "/openshell.compute.v1.ComputeDriver/AttackerControlled12345" + ), + ("driver.unknown", "unknown"), + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn podman_driver_spans_reach_otlp_collector_with_distinct_service_name() { + let _tracing_lock = super::test_lock().await; + let received = Arc::new(Mutex::new(Received::default())); + let exported = Arc::new(tokio::sync::Notify::new()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let address = listener.local_addr().unwrap(); + let collector = Collector { + received: Arc::clone(&received), + exported: Arc::clone(&exported), + }; + let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel(); + let server = tokio::spawn(async move { + tonic::transport::Server::builder() + .add_service(TraceServiceServer::new(collector)) + .serve_with_incoming_shutdown( + tokio_stream::wrappers::TcpListenerStream::new(listener), + async { + let _ = shutdown_rx.await; + }, + ) + .await + }); + + let (provider, error) = super::provider_for(Some(&format!("http://{address}"))); + assert!(error.is_none()); + let provider = provider.expect("provider"); + let subscriber = tracing_subscriber::registry().with(super::layer(&provider)); + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!("podman.create_sandbox", sandbox.id = "sb-otlp"); + drop(span.enter()); + drop(span); + }); + let export_completed = exported.notified(); + provider.force_flush().unwrap(); + tokio::time::timeout(std::time::Duration::from_secs(5), export_completed) + .await + .expect("OTLP export should complete"); + provider.shutdown().unwrap(); + shutdown_tx.send(()).unwrap(); + server.await.unwrap().unwrap(); + + let received = received.lock().unwrap(); + assert!( + received + .spans + .iter() + .any(|span| span.name == "podman.create_sandbox") + ); + assert!( + received + .service_names + .iter() + .any(|name| name == "openshell-driver-podman") + ); + } +} diff --git a/crates/openshell-driver-vm/src/otel_tracing.rs b/crates/openshell-driver-vm/src/otel_tracing.rs index adfeb896c6..7ae75357db 100644 --- a/crates/openshell-driver-vm/src/otel_tracing.rs +++ b/crates/openshell-driver-vm/src/otel_tracing.rs @@ -5,8 +5,8 @@ use http::Request; use openshell_otel::{ - HeaderMapExtractor, OtlpTraceConfig, RecordGrpcFailure, SdkTracerProvider, ServiceName, - SetupError, + HeaderMapExtractor, OtlpTraceConfig, RecordGrpcFailure, RecordGrpcStatus, SdkTracerProvider, + ServiceName, SetupError, }; use opentelemetry::propagation::TextMapPropagator; use opentelemetry::trace::TraceContextExt as _; @@ -21,14 +21,21 @@ const INSTRUMENTATION_SCOPE: &str = "openshell-driver-vm"; const COMPUTE_DRIVER_SERVICE: &str = "openshell.compute.v1.ComputeDriver"; /// Trace every inbound compute-driver RPC at the tonic service boundary. -pub fn compute_driver_rpc_layer() --> TraceLayer { +pub fn compute_driver_rpc_layer() -> TraceLayer< + GrpcMakeClassifier, + ComputeDriverRpcSpan, + (), + RecordGrpcStatus, + (), + RecordGrpcStatus, + RecordGrpcFailure, +> { TraceLayer::new_for_grpc() .make_span_with(ComputeDriverRpcSpan) .on_request(()) - .on_response(()) + .on_response(RecordGrpcStatus) .on_body_chunk(()) - .on_eos(()) + .on_eos(RecordGrpcStatus) .on_failure(RecordGrpcFailure) } diff --git a/crates/openshell-otel/Cargo.toml b/crates/openshell-otel/Cargo.toml index b53a716819..8155989824 100644 --- a/crates/openshell-otel/Cargo.toml +++ b/crates/openshell-otel/Cargo.toml @@ -23,6 +23,7 @@ tonic = { workspace = true } tower-http = { workspace = true } [dev-dependencies] +opentelemetry_sdk = { workspace = true, features = ["testing"] } tokio = { workspace = true } [lints] diff --git a/crates/openshell-otel/src/grpc.rs b/crates/openshell-otel/src/grpc.rs index 9eb7caa488..564f889aeb 100644 --- a/crates/openshell-otel/src/grpc.rs +++ b/crates/openshell-otel/src/grpc.rs @@ -4,7 +4,7 @@ //! Shared gRPC tracing adapters. use tower_http::classify::GrpcFailureClass; -use tower_http::trace::OnFailure; +use tower_http::trace::{OnEos, OnFailure, OnResponse}; use tracing::Span; /// Records a non-OK gRPC outcome on the request span. @@ -24,3 +24,118 @@ impl OnFailure for RecordGrpcFailure { } } } + +/// Records a gRPC status from response headers or trailers. +#[derive(Debug, Clone, Copy)] +pub struct RecordGrpcStatus; + +impl RecordGrpcStatus { + fn record(headers: &http::HeaderMap, span: &Span) { + let Some(code) = headers + .get("grpc-status") + .and_then(|status| status.to_str().ok()) + .and_then(|status| status.parse::().ok()) + else { + return; + }; + if code != tonic::Code::Ok as i32 { + crate::mark_error(span); + } + span.record("rpc.grpc.status_code", code); + } +} + +impl OnResponse for RecordGrpcStatus { + fn on_response(self, response: &http::Response, _latency: std::time::Duration, span: &Span) { + Self::record(response.headers(), span); + } +} + +impl OnEos for RecordGrpcStatus { + fn on_eos( + self, + trailers: Option<&http::HeaderMap>, + _stream_duration: std::time::Duration, + span: &Span, + ) { + if let Some(trailers) = trailers { + Self::record(trailers, span); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use opentelemetry_sdk::trace::{InMemorySpanExporterBuilder, SdkTracerProvider}; + use tracing_subscriber::layer::SubscriberExt as _; + + #[test] + fn grpc_status_records_and_marks_non_ok_trailer_status() { + let _tracing_lock = crate::test_lock(); + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::layer(&provider, "test")); + let mut trailers = http::HeaderMap::new(); + trailers.insert("grpc-status", http::HeaderValue::from_static("13")); + + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!( + "rpc", + otel.status_code = tracing::field::Empty, + rpc.grpc.status_code = tracing::field::Empty, + ); + RecordGrpcStatus.on_eos(Some(&trailers), std::time::Duration::ZERO, &span); + }); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + assert_eq!(spans.len(), 1); + assert!(matches!( + spans[0].status, + opentelemetry::trace::Status::Error { .. } + )); + assert!(spans[0].attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" && attribute.value.to_string() == "13" + })); + provider.shutdown().unwrap(); + } + + #[test] + fn grpc_status_records_header_status_without_eos_overwrite() { + let _tracing_lock = crate::test_lock(); + let exporter = InMemorySpanExporterBuilder::new().build(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with(crate::layer(&provider, "test")); + let response = http::Response::builder() + .header("grpc-status", "13") + .body(()) + .unwrap(); + + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!( + "rpc", + otel.status_code = tracing::field::Empty, + rpc.grpc.status_code = tracing::field::Empty, + ); + RecordGrpcStatus.on_response(&response, std::time::Duration::ZERO, &span); + RecordGrpcStatus.on_eos(None, std::time::Duration::ZERO, &span); + }); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + assert_eq!(spans.len(), 1); + assert!(matches!( + spans[0].status, + opentelemetry::trace::Status::Error { .. } + )); + assert!(spans[0].attributes.iter().any(|attribute| { + attribute.key.as_str() == "rpc.grpc.status_code" && attribute.value.to_string() == "13" + })); + provider.shutdown().unwrap(); + } +} diff --git a/crates/openshell-otel/src/lib.rs b/crates/openshell-otel/src/lib.rs index 1536f66f24..3865670b0b 100644 --- a/crates/openshell-otel/src/lib.rs +++ b/crates/openshell-otel/src/lib.rs @@ -6,7 +6,7 @@ mod grpc; mod propagation; -pub use grpc::RecordGrpcFailure; +pub use grpc::{RecordGrpcFailure, RecordGrpcStatus}; pub use propagation::{HeaderMapExtractor, MetadataMapInjector, TraceContextInterceptor}; use opentelemetry::KeyValue; @@ -18,6 +18,7 @@ pub use opentelemetry_sdk::trace::SdkTracerProvider; use tracing::Subscriber; use tracing_opentelemetry::OpenTelemetryLayer; use tracing_subscriber::Layer as _; +use tracing_subscriber::layer::{Context, Filter}; use tracing_subscriber::registry::LookupSpan; const SDK_UNKNOWN_SERVICE_PREFIX: &str = "unknown_service"; @@ -193,6 +194,23 @@ pub type OtlpLayer = tracing_subscriber::filter::Filtered< S, >; +pub type TargetOtlpLayer = + tracing_subscriber::filter::Filtered, TargetPrefixFilter, S>; + +#[derive(Debug, Clone, Copy)] +pub struct TargetPrefixFilter { + prefix: &'static str, + include: bool, +} + +impl Filter for TargetPrefixFilter { + fn enabled(&self, metadata: &tracing::Metadata<'_>, _ctx: &Context<'_, S>) -> bool { + metadata.is_span() + && !metadata.target().starts_with("opentelemetry") + && (metadata.target().starts_with(self.prefix) == self.include) + } +} + /// Build a tracing layer that exports spans and excludes exporter callsites. pub fn layer(provider: &SdkTracerProvider, instrumentation_scope: &'static str) -> OtlpLayer where @@ -205,6 +223,47 @@ where })) } +/// Build a tracing layer that exports only spans from `target_prefix`. +pub fn layer_for_target_prefix( + provider: &SdkTracerProvider, + instrumentation_scope: &'static str, + target_prefix: &'static str, +) -> TargetOtlpLayer +where + S: Subscriber + for<'span> LookupSpan<'span>, +{ + tracing_opentelemetry::layer() + .with_tracer(provider.tracer(instrumentation_scope)) + .with_filter(TargetPrefixFilter { + prefix: target_prefix, + include: true, + }) +} + +/// Build a tracing layer that excludes spans from `target_prefix`. +pub fn layer_excluding_target_prefix( + provider: &SdkTracerProvider, + instrumentation_scope: &'static str, + target_prefix: &'static str, +) -> TargetOtlpLayer +where + S: Subscriber + for<'span> LookupSpan<'span>, +{ + tracing_opentelemetry::layer() + .with_tracer(provider.tracer(instrumentation_scope)) + .with_filter(TargetPrefixFilter { + prefix: target_prefix, + include: false, + }) +} + +#[cfg(test)] +pub(crate) fn test_lock() -> std::sync::MutexGuard<'static, ()> { + static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(()); + LOCK.lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) +} + #[cfg(test)] mod tests { use super::*; @@ -213,6 +272,7 @@ mod tests { fn error_status_guard_marks_only_unfinished_results() { use tracing_subscriber::layer::SubscriberExt as _; + let _tracing_lock = test_lock(); let exporter = opentelemetry_sdk::trace::InMemorySpanExporterBuilder::new().build(); let provider = SdkTracerProvider::builder() .with_simple_exporter(exporter.clone()) @@ -255,6 +315,65 @@ mod tests { assert_eq!(succeeded.status, opentelemetry::trace::Status::Unset); } + #[test] + fn target_scoped_layers_partition_in_process_driver_spans() { + use tracing_subscriber::layer::SubscriberExt as _; + + let _tracing_lock = test_lock(); + let gateway_exporter = opentelemetry_sdk::trace::InMemorySpanExporterBuilder::new().build(); + let gateway_provider = SdkTracerProvider::builder() + .with_simple_exporter(gateway_exporter.clone()) + .build(); + let driver_exporter = opentelemetry_sdk::trace::InMemorySpanExporterBuilder::new().build(); + let driver_provider = SdkTracerProvider::builder() + .with_simple_exporter(driver_exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry() + .with(layer_excluding_target_prefix( + &gateway_provider, + "gateway", + "driver_target", + )) + .with(layer_for_target_prefix( + &driver_provider, + "driver", + "driver_target", + )); + + tracing::subscriber::with_default(subscriber, || { + let gateway = tracing::info_span!(target: "gateway_target", "gateway.span"); + let _entered = gateway.enter(); + drop(tracing::info_span!(target: "driver_target::operation", "driver.span")); + }); + + gateway_provider.force_flush().unwrap(); + driver_provider.force_flush().unwrap(); + let gateway_spans = gateway_exporter.get_finished_spans().unwrap(); + let driver_spans = driver_exporter.get_finished_spans().unwrap(); + assert_eq!( + gateway_spans + .iter() + .map(|span| span.name.as_ref()) + .collect::>(), + ["gateway.span"] + ); + assert_eq!( + driver_spans + .iter() + .map(|span| span.name.as_ref()) + .collect::>(), + ["driver.span"] + ); + assert_eq!( + driver_spans[0].span_context.trace_id(), + gateway_spans[0].span_context.trace_id() + ); + assert_eq!( + driver_spans[0].parent_span_id, + gateway_spans[0].span_context.span_id() + ); + } + #[test] fn resource_uses_fixed_service_identity_and_custom_attributes() { let resource = resource_for(&OtlpTraceConfig { diff --git a/crates/openshell-server/src/cli.rs b/crates/openshell-server/src/cli.rs index 2e86c3a1b5..7d9d9cac47 100644 --- a/crates/openshell-server/src/cli.rs +++ b/crates/openshell-server/src/cli.rs @@ -17,7 +17,10 @@ use crate::certgen; use crate::compute::driver_config::GuestTlsPaths; use crate::config_file::{self, ConfigFile, GatewayFileSection}; use crate::defaults::{self, LocalTlsPaths}; -use crate::{ServerStartupConfig, run_server, tracing_bus::TracingLogBus}; +use crate::{ + ServerStartupConfig, configured_compute_driver_for_startup, run_server, + tracing_bus::TracingLogBus, +}; /// `OpenShell` gateway process - gRPC and HTTP server with protocol multiplexing. /// @@ -470,6 +473,7 @@ fn prepare_server_config(args: &mut RunArgs, matches: &ArgMatches) -> Result Result<()> { let prepared = prepare_server_config(&mut args, &matches)?; + let compute_driver = configured_compute_driver_for_startup(&prepared)?; let tracing_log_bus = TracingLogBus::new(); let otlp_config = prepared @@ -481,6 +485,7 @@ async fn run_from_args(mut args: RunArgs, matches: ArgMatches) -> Result<()> { .unwrap_or_else(|_| EnvFilter::new(&prepared.config.log_level)), &tracing_log_bus, otlp_config, + crate::tracing_setup::podman_export_enabled(&compute_driver), ); let has_client_ca = prepared @@ -537,7 +542,7 @@ async fn run_from_args(mut args: RunArgs, matches: ArgMatches) -> Result<()> { info!(bind = %prepared.config.bind_address, "Starting OpenShell server"); - let result = Box::pin(run_server(prepared, tracing_log_bus)).await; + let result = Box::pin(run_server(prepared, compute_driver, tracing_log_bus)).await; tracing_handle.shutdown(); diff --git a/crates/openshell-server/src/compute/mod.rs b/crates/openshell-server/src/compute/mod.rs index 82bf2e2c58..0dce413dc1 100644 --- a/crates/openshell-server/src/compute/mod.rs +++ b/crates/openshell-server/src/compute/mod.rs @@ -800,7 +800,7 @@ impl ComputeRuntime { let driver = PodmanComputeDriver::new(config) .await .map_err(|err| ComputeError::Message(err.to_string()))?; - let driver: SharedComputeDriver = Arc::new(PodmanDriverService::new(driver)); + let driver: SharedComputeDriver = Arc::new(PodmanDriverService::new_in_process(driver)); Self::from_driver( ComputeDriverKind::Podman.as_str().to_string(), driver, diff --git a/crates/openshell-server/src/lib.rs b/crates/openshell-server/src/lib.rs index 979e372084..86c8e28e2b 100644 --- a/crates/openshell-server/src/lib.rs +++ b/crates/openshell-server/src/lib.rs @@ -437,6 +437,7 @@ impl ServerState { /// Returns an error if the server fails to start or encounters a fatal error. pub(crate) async fn run_server( startup: ServerStartupConfig, + compute_driver: ConfiguredComputeDriver, tracing_log_bus: TracingLogBus, ) -> Result<()> { let ServerStartupConfig { @@ -593,6 +594,7 @@ pub(crate) async fn run_server( let (compute, operator_allowlist) = build_compute_runtime( &config, driver_startup, + compute_driver, store.clone(), sandbox_index.clone(), sandbox_watch_bus.clone(), @@ -1088,6 +1090,7 @@ type OperatorAllowlistArc = Option, + driver: ConfiguredComputeDriver, store: Arc, sandbox_index: SandboxIndex, sandbox_watch_bus: SandboxWatchBus, @@ -1095,7 +1098,6 @@ async fn build_compute_runtime( supervisor_sessions: Arc, shutdown_rx: watch::Receiver, ) -> Result<(ComputeRuntime, OperatorAllowlistArc)> { - let driver = configured_compute_driver(config, driver_startup)?; info!(driver = %driver.name(), "Using compute driver"); let (runtime, operator_allowlist) = match driver { @@ -1205,7 +1207,7 @@ async fn build_compute_runtime( } #[derive(Debug, Clone)] -enum ConfiguredComputeDriver { +pub(crate) enum ConfiguredComputeDriver { Builtin(ComputeDriverKind), Remote { name: String }, } @@ -1242,6 +1244,21 @@ fn configured_compute_driver( } } +pub(crate) fn configured_compute_driver_for_startup( + startup: &ServerStartupConfig, +) -> Result { + configured_compute_driver( + &startup.config, + compute::driver_config::DriverStartupContext { + file: startup.config_file.as_ref(), + guest_tls: startup.guest_tls.as_ref(), + gateway_port: startup.config.bind_address.port(), + gateway_tls_enabled: startup.config.tls.is_some(), + endpoint_overrides: &startup.config.compute_driver_endpoints, + }, + ) +} + fn resolve_configured_compute_driver( driver_name: &str, driver_startup: compute::driver_config::DriverStartupContext<'_>, diff --git a/crates/openshell-server/src/otel_tracing.rs b/crates/openshell-server/src/otel_tracing.rs index cbe23f4231..d29f3b77f8 100644 --- a/crates/openshell-server/src/otel_tracing.rs +++ b/crates/openshell-server/src/otel_tracing.rs @@ -90,11 +90,15 @@ pub fn provider_for(cfg: Option<&OtlpConfig>) -> (Option, Opt /// /// Events stay on the gateway's logging layers. Spans emitted by the /// OpenTelemetry crates are excluded to prevent recursive export traffic. -pub fn layer(provider: &SdkTracerProvider) -> openshell_otel::OtlpLayer +pub fn layer(provider: &SdkTracerProvider) -> openshell_otel::TargetOtlpLayer where S: Subscriber + for<'span> LookupSpan<'span>, { - openshell_otel::layer(provider, INSTRUMENTATION_SCOPE) + openshell_otel::layer_excluding_target_prefix( + provider, + INSTRUMENTATION_SCOPE, + openshell_driver_podman::otel_tracing::IN_PROCESS_TARGET_PREFIX, + ) } /// Isolated in-memory span exporters for tracing tests. diff --git a/crates/openshell-server/src/tracing_setup.rs b/crates/openshell-server/src/tracing_setup.rs index 321edefafe..ac3c1ae79b 100644 --- a/crates/openshell-server/src/tracing_setup.rs +++ b/crates/openshell-server/src/tracing_setup.rs @@ -11,12 +11,14 @@ use opentelemetry_sdk::trace::SdkTracerProvider; use tracing_subscriber::EnvFilter; use tracing_subscriber::prelude::*; +use crate::ConfiguredComputeDriver; use crate::config_file::OtlpConfig; use crate::otel_tracing::SetupError; use crate::tracing_bus::TracingLogBus; pub struct TracingHandle { tracer_provider: Option, + podman_tracer_provider: Option, } impl TracingHandle { @@ -26,22 +28,77 @@ impl TracingHandle { { tracing::warn!(error = %err, "OTLP tracer provider shutdown failed"); } + if let Some(provider) = &self.podman_tracer_provider + && let Err(err) = provider.shutdown() + { + tracing::warn!(error = %err, "Podman OTLP tracer provider shutdown failed"); + } } } +#[must_use] +pub fn podman_export_enabled(driver: &ConfiguredComputeDriver) -> bool { + matches!( + driver, + ConfiguredComputeDriver::Builtin(openshell_core::ComputeDriverKind::Podman) + ) +} + pub fn install( env_filter: EnvFilter, tracing_log_bus: &TracingLogBus, otlp_config: Option<&OtlpConfig>, + enable_podman_export: bool, ) -> (TracingHandle, Option) { let (tracer_provider, setup_error) = crate::otel_tracing::provider_for(otlp_config); + let podman_endpoint = enable_podman_export + .then_some(otlp_config) + .flatten() + .map(|config| config.endpoint.as_str()); + let (podman_tracer_provider, podman_setup_error) = + openshell_driver_podman::otel_tracing::provider_for(podman_endpoint); tracing_subscriber::registry() .with(env_filter) .with(tracing_subscriber::fmt::layer()) .with(tracing_log_bus.layer()) .with(tracer_provider.as_ref().map(crate::otel_tracing::layer)) + .with( + podman_tracer_provider + .as_ref() + .map(openshell_driver_podman::otel_tracing::in_process_layer), + ) .init(); - (TracingHandle { tracer_provider }, setup_error) + ( + TracingHandle { + tracer_provider, + podman_tracer_provider, + }, + setup_error.or(podman_setup_error), + ) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn podman_export_is_enabled_only_when_podman_is_selected() { + use crate::ConfiguredComputeDriver; + use openshell_core::ComputeDriverKind; + + assert!(podman_export_enabled(&ConfiguredComputeDriver::Builtin( + ComputeDriverKind::Podman + ))); + assert!(!podman_export_enabled(&ConfiguredComputeDriver::Builtin( + ComputeDriverKind::Docker + ))); + assert!(!podman_export_enabled(&ConfiguredComputeDriver::Builtin( + ComputeDriverKind::Kubernetes + ))); + assert!(!podman_export_enabled(&ConfiguredComputeDriver::Remote { + name: "custom".to_string(), + })); + } }