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/bazel-OpenShell b/bazel-OpenShell new file mode 120000 index 0000000000..c12b74a246 --- /dev/null +++ b/bazel-OpenShell @@ -0,0 +1 @@ +/Users/khicks/Library/Caches/bazel/_bazel_khicks/7c8977f45119661e0b30646169fa9a1e/execroot/_main \ No newline at end of file diff --git a/bazel-bin b/bazel-bin new file mode 120000 index 0000000000..c8068eed74 --- /dev/null +++ b/bazel-bin @@ -0,0 +1 @@ +/Users/khicks/Library/Caches/bazel/_bazel_khicks/7c8977f45119661e0b30646169fa9a1e/execroot/_main/bazel-out/darwin_arm64-fastbuild/bin \ No newline at end of file diff --git a/bazel-out b/bazel-out new file mode 120000 index 0000000000..4d71984392 --- /dev/null +++ b/bazel-out @@ -0,0 +1 @@ +/Users/khicks/Library/Caches/bazel/_bazel_khicks/7c8977f45119661e0b30646169fa9a1e/execroot/_main/bazel-out \ No newline at end of file diff --git a/bazel-testlogs b/bazel-testlogs new file mode 120000 index 0000000000..685bce1d9d --- /dev/null +++ b/bazel-testlogs @@ -0,0 +1 @@ +/Users/khicks/Library/Caches/bazel/_bazel_khicks/7c8977f45119661e0b30646169fa9a1e/execroot/_main/bazel-out/darwin_arm64-fastbuild/testlogs \ No newline at end of file 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(), + })); + } }