From 215fa4510c456bf555577bfebeec756e1dd35838 Mon Sep 17 00:00:00 2001 From: Kris Hicks Date: Wed, 29 Jul 2026 16:25:30 -0700 Subject: [PATCH] feat(vm): export driver traces over OTLP Continue distributed traces across the gateway-to-driver process boundary and export VM driver spans to the same OTLP/gRPC collector. The driver reports as the distinct openshell-driver-vm service. Updated the gateway architecture and configuration reference with a generic external-driver forwarding contract. Instrumented: - Every RemoteComputeDriver RPC injects the active W3C trace context into tonic metadata. Managed VM readiness and runtime initialization give startup capability probes stable parent operations rather than isolated root spans. - A tonic service layer creates fixed, low-cardinality server spans for every ComputeDriver RPC. New handlers inherit tracing automatically; failures record OpenTelemetry error status and the gRPC status code. - Background provisioning remains attached to CreateSandbox after the RPC returns without extending the RPC span lifetime. - Provisioning records image preparation, bootstrap image resolution, overlay preparation, lifecycle configuration, pre-launch hooks, guest preparation, and launcher spawn as child spans. - VM startup reconciliation roots one trace for the persisted-sandbox scan, with per-sandbox restore and provision operations beneath it. The root remains open until all spawned restore tasks finish. - Delete cleanup records its own child operation. Design notes: - The gateway forwards its configured OTLP endpoint to managed external drivers. SDK `OTEL_*` variables continue to own sampling, batching, limits, headers, and transport tuning. - The VM driver has its own tracer provider and service resource so trace backends preserve the service boundary. - RPC operation names come from an explicit method mapping, keeping cardinality bounded without parsing the protobuf descriptor set at runtime. - Propagation uses a remote SpanContext for spawned provisioning. This keeps one trace while allowing the CreateSandbox server span to finish when the RPC response is sent. - Startup restoration is independent of gateway requests. It begins at the VM driver reconciliation span rather than attaching to an unrelated RPC. - Existing tracing events remain on the logging path. The OpenTelemetry layer exports spans only and excludes the SDK exporter callsites to avoid recursive traces. - Export configuration failures do not prevent the driver from serving, and buffered spans are drained during graceful shutdown. - Trace fields identify drivers, sandboxes, images, lifecycle phases, and gRPC outcomes without recording credentials, sandbox tokens, or request query parameters. Refs #2507 Signed-off-by: Kris Hicks --- Cargo.lock | 8 + architecture/gateway.md | 21 +- crates/openshell-driver-vm/Cargo.toml | 9 + crates/openshell-driver-vm/src/driver.rs | 648 +++++++++++++++++- crates/openshell-driver-vm/src/lib.rs | 1 + crates/openshell-driver-vm/src/lifecycle.rs | 18 + crates/openshell-driver-vm/src/main.rs | 45 +- .../openshell-driver-vm/src/otel_tracing.rs | 242 +++++++ crates/openshell-server/src/compute/mod.rs | 143 +++- crates/openshell-server/src/compute/vm.rs | 88 ++- crates/openshell-server/src/lib.rs | 5 +- crates/openshell-server/src/otel_tracing.rs | 33 + crates/openshell-server/src/test_support.rs | 29 +- docs/reference/gateway-config.mdx | 4 +- 14 files changed, 1240 insertions(+), 54 deletions(-) create mode 100644 crates/openshell-driver-vm/src/otel_tracing.rs diff --git a/Cargo.lock b/Cargo.lock index 528309a970..3026807275 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3968,14 +3968,19 @@ dependencies = [ "clap", "flate2", "futures", + "http 1.4.0", "libc", "libloading", "miette", "nix 0.29.0", "oci-client", "openshell-core", + "openshell-otel", "openshell-policy", "openshell-vfio", + "opentelemetry", + "opentelemetry-proto", + "opentelemetry_sdk", "polling", "prost", "prost-types", @@ -3985,10 +3990,13 @@ dependencies = [ "sha2 0.10.9", "tar", "temp-env", + "tempfile", "tokio", "tokio-stream", "tonic", + "tower-http 0.6.8", "tracing", + "tracing-opentelemetry", "tracing-subscriber", "url", "zstd", diff --git a/architecture/gateway.md b/architecture/gateway.md index c7569c3d5b..602426064b 100644 --- a/architecture/gateway.md +++ b/architecture/gateway.md @@ -639,15 +639,18 @@ to export; the SDK's `OTEL_*` variables tune how. Transport is OTLP over gRPC only. Shared provider, resource, and tracing-layer construction lives in `openshell-otel`. -Span emission requires no per-handler instrumentation. The `tower_http` -`TraceLayer` in `multiplex.rs` opens a span per inbound request, and that span -continues incoming W3C trace context when present or starts a new trace -otherwise. It is named for the RPC and carries the request ID that also appears -in the gateway's logs — the identifier that lets an operator pivot between a -trace and its log lines. Store and compute-driver spans become children of the -request span. Reconciliation, provider refresh, and driver-watch loops create -their own operation spans because they have no inbound request to provide a -parent. gRPC status is recorded when response trailers arrive. +The `tower_http` `TraceLayer` in `multiplex.rs` opens a span per inbound request, +and that span continues incoming W3C trace context when present or starts a new +trace otherwise. It is named for the RPC and carries the request ID that also +appears in the gateway's logs — the identifier that lets an operator pivot +between a trace and its log lines. Store and compute-driver spans become +children of the request span. Reconciliation, provider refresh, and +driver-watch loops create their own operation spans because they have no +inbound request to provide a parent. gRPC status is recorded when response +trailers arrive. + +The gateway forwards OTLP configuration and W3C trace context to managed +external drivers. Each driver exports under its own service name. Two invariants shape the failure behavior. Telemetry is diagnostic, so no OTLP failure stops the gateway from serving: a malformed endpoint is logged at diff --git a/crates/openshell-driver-vm/Cargo.toml b/crates/openshell-driver-vm/Cargo.toml index cef3e67f88..0ee2740c77 100644 --- a/crates/openshell-driver-vm/Cargo.toml +++ b/crates/openshell-driver-vm/Cargo.toml @@ -20,12 +20,15 @@ path = "src/main.rs" [dependencies] openshell-core = { path = "../openshell-core", default-features = false } +openshell-otel = { path = "../openshell-otel" } openshell-policy = { path = "../openshell-policy" } openshell-vfio = { path = "../openshell-vfio" } bollard = { version = "0.20", features = ["ssh"] } tokio = { workspace = true } tonic = { workspace = true, features = ["transport"] } +tower-http = { workspace = true } +http = { workspace = true } prost = { workspace = true } prost-types = { workspace = true } futures = { workspace = true } @@ -34,6 +37,9 @@ nix = { workspace = true } clap = { workspace = true } tracing = { workspace = true } tracing-subscriber = { workspace = true } +opentelemetry = { workspace = true } +opentelemetry_sdk = { workspace = true } +tracing-opentelemetry = { workspace = true } miette = { workspace = true } url = { workspace = true } serde = { workspace = true } @@ -56,6 +62,9 @@ telemetry = ["openshell-core/telemetry"] [dev-dependencies] temp-env = "0.3" +tempfile = "3" +opentelemetry_sdk = { workspace = true, features = ["testing"] } +opentelemetry-proto = { version = "0.32", default-features = false, features = ["gen-tonic", "trace"] } # smol-rs/polling drives the BSD/macOS parent-death detection in # procguard via kqueue's EVFILT_PROC / NOTE_EXIT filter. We could use diff --git a/crates/openshell-driver-vm/src/driver.rs b/crates/openshell-driver-vm/src/driver.rs index 841457f38c..1fd275310d 100644 --- a/crates/openshell-driver-vm/src/driver.rs +++ b/crates/openshell-driver-vm/src/driver.rs @@ -52,6 +52,7 @@ use openshell_core::proto_struct::{ deserialize_optional_non_empty_string_list, struct_to_json_value, }; use openshell_vfio::SysfsRoot; +use opentelemetry::trace::TraceContextExt as _; use prost::Message; use sha2::{Digest, Sha256}; use std::collections::{HashMap, HashSet}; @@ -72,7 +73,8 @@ use tokio::sync::{Mutex, broadcast, mpsc}; use tokio::task::JoinHandle; use tokio_stream::wrappers::ReceiverStream; use tonic::{Request, Response, Status}; -use tracing::{info, warn}; +use tracing::{Instrument as _, info, warn}; +use tracing_opentelemetry::OpenTelemetrySpanExt as _; use url::{Host, Url}; const DRIVER_NAME: &str = "openshell-driver-vm"; @@ -392,6 +394,27 @@ enum OverlayPreparation { PreserveExisting, } +fn provisioning_span( + parent: &opentelemetry::Context, + sandbox_id: &str, + image_ref: &str, +) -> tracing::Span { + let span = tracing::info_span!( + parent: None, + "vm.provision", + otel.name = "vm.provision", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + image.ref = %image_ref, + ); + let parent_span_context = parent.span().span_context().clone(); + if parent_span_context.is_valid() { + let parent = opentelemetry::Context::new().with_remote_span_context(parent_span_context); + let _ = span.set_parent(parent); + } + span +} + #[derive(Clone)] pub struct VmDriver { config: VmDriverConfig, @@ -600,17 +623,22 @@ impl VmDriver { let sandbox_id = sandbox.id.clone(); let image_ref_for_task = image_ref.clone(); let state_dir_for_task = state_dir.clone(); - let task = tokio::spawn(async move { - driver - .provision_sandbox( - sandbox_for_task, - image_ref_for_task, - state_dir_for_task, - tls_paths, - OverlayPreparation::Fresh, - ) - .await; - }); + let parent = tracing::Span::current().context(); + let provisioning_span = provisioning_span(&parent, &sandbox_id, &image_ref); + let task = tokio::spawn( + async move { + driver + .provision_sandbox( + sandbox_for_task, + image_ref_for_task, + state_dir_for_task, + tls_paths, + OverlayPreparation::Fresh, + ) + .await; + } + .instrument(provisioning_span), + ); let mut registry = self.registry.lock().await; if let Some(record) = registry.get_mut(&sandbox_id) { @@ -645,6 +673,7 @@ impl VmDriver { ) .await { + tracing::Span::current().record("otel.status_code", "ERROR"); if err.code() == tonic::Code::Cancelled { if overlay_preparation == OverlayPreparation::Fresh { let _ = tokio::fs::remove_dir_all(&state_dir).await; @@ -942,7 +971,7 @@ impl VmDriver { console_output = %console_output.display(), "vm driver: spawning VM launcher" ); - let child = match command.spawn() { + let child = match spawn_vm_launcher(&mut command, &sandbox.id, &plan.backend) { Ok(child) => child, Err(err) => { warn!( @@ -1029,6 +1058,16 @@ impl VmDriver { Ok(()) } + #[tracing::instrument( + name = "vm.delete", + skip(self), + fields( + otel.name = "vm.delete", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + sandbox.name = %sandbox_name, + ) + )] pub async fn delete_sandbox( &self, sandbox_id: &str, @@ -1146,6 +1185,14 @@ impl VmDriver { snapshots } + #[tracing::instrument( + name = "reconcile", + skip_all, + fields( + otel.name = "reconcile.sandboxes", + driver.name = "vm", + ) + )] async fn restore_persisted_sandboxes(&self) { let state_root = sandboxes_root_dir(&self.config.state_dir); let mut entries = match tokio::fs::read_dir(&state_root).await { @@ -1216,11 +1263,17 @@ impl VmDriver { continue; } - self.restore_persisted_sandbox(sandbox, state_dir).await; + self.restore_persisted_sandbox(sandbox, state_dir, &tracing::Span::current()) + .await; } } - async fn restore_persisted_sandbox(&self, sandbox: Sandbox, state_dir: PathBuf) { + async fn restore_persisted_sandbox( + &self, + sandbox: Sandbox, + state_dir: PathBuf, + reconciliation_span: &tracing::Span, + ) { let Some(image_ref) = self.resolved_sandbox_image(&sandbox) else { warn!( sandbox_id = %sandbox.id, @@ -1301,17 +1354,32 @@ impl VmDriver { let driver = self.clone(); let sandbox_id = sandbox.id.clone(); - let task = tokio::spawn(async move { - driver - .provision_sandbox( - sandbox, - image_ref, - state_dir, - tls_paths, - OverlayPreparation::PreserveExisting, - ) - .await; - }); + let restoration_span = tracing::info_span!( + parent: reconciliation_span, + "vm.restore", + otel.name = "vm.restore", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + ); + let reconciliation_span = reconciliation_span.clone(); + let provisioning_span = + provisioning_span(&restoration_span.context(), &sandbox_id, &image_ref); + let task = tokio::spawn( + async move { + driver + .provision_sandbox( + sandbox, + image_ref, + state_dir, + tls_paths, + OverlayPreparation::PreserveExisting, + ) + .await; + drop(reconciliation_span); + } + .instrument(provisioning_span) + .instrument(restoration_span), + ); let mut registry = self.registry.lock().await; if let Some(record) = registry.get_mut(&sandbox_id) { @@ -1686,6 +1754,16 @@ impl VmDriver { } } + #[tracing::instrument( + name = "vm.prepare_images", + skip(self), + fields( + otel.name = "vm.prepare_images", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + image.ref = %image_ref, + ) + )] async fn prepare_runtime_images( &self, sandbox_id: &str, @@ -1729,6 +1807,16 @@ impl VmDriver { sandbox_image_ref.to_string() } + #[tracing::instrument( + name = "vm.prepare_overlay", + skip_all, + fields( + otel.name = "vm.prepare_overlay", + otel.status_code = tracing::field::Empty, + overlay.path = %overlay_disk.display(), + preparation = ?preparation, + ) + )] async fn prepare_runtime_overlay( &self, overlay_disk: &Path, @@ -1787,6 +1875,16 @@ impl VmDriver { }) } + #[tracing::instrument( + name = "vm.resolve_bootstrap_image", + skip(self), + fields( + otel.name = "vm.resolve_bootstrap_image", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + image.ref = %image_ref, + ) + )] async fn ensure_cached_bootstrap_rootfs_image( &self, sandbox_id: &str, @@ -3109,7 +3207,11 @@ impl ComputeDriver for VmDriver { } loop { - match rx.recv().await { + let event = tokio::select! { + () = tx.closed() => return, + event = rx.recv() => event, + }; + match event { Ok(event) => { if let Some(watch_sandboxes_event::Payload::Sandbox(sandbox_event)) = &event.payload @@ -3128,7 +3230,8 @@ impl ComputeDriver for VmDriver { } }); - Ok(Response::new(Box::pin(ReceiverStream::new(out_rx)))) + let stream: Self::WatchSandboxesStream = Box::pin(ReceiverStream::new(out_rx)); + Ok(Response::new(stream)) } } @@ -4677,6 +4780,16 @@ fn inject_guest_sandbox_token(overlay_disk: &Path, token: &str) -> Result<(), St } #[allow(clippy::result_large_err)] +#[tracing::instrument( + name = "vm.prepare_guest", + skip(dropins), + fields( + otel.name = "vm.prepare_guest", + otel.status_code = tracing::field::Empty, + overlay.path = %overlay_disk.display(), + dropin.count = dropins.len(), + ) +)] fn inject_guest_init_dropins( overlay_disk: &Path, dropins: &[GuestInitDropin], @@ -4979,6 +5092,24 @@ async fn terminate_vm_process(child: &mut Child) -> Result<(), std::io::Error> { } } +#[tracing::instrument( + name = "vm.launch", + skip(command), + fields( + otel.name = "vm.launch", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox_id, + vm.backend = ?backend, + ) +)] +fn spawn_vm_launcher( + command: &mut Command, + sandbox_id: &str, + backend: &VmBackend, +) -> Result { + command.spawn() +} + fn sandbox_snapshot(sandbox: &Sandbox, condition: SandboxCondition, deleting: bool) -> Sandbox { Sandbox { id: sandbox.id.clone(), @@ -5161,6 +5292,467 @@ mod tests { static ENV_LOCK: std::sync::LazyLock> = std::sync::LazyLock::new(|| std::sync::Mutex::new(())); + static TRACE_EXPORTER: std::sync::LazyLock = + std::sync::LazyLock::new(|| { + use opentelemetry::trace::TracerProvider as _; + use tracing_subscriber::layer::SubscriberExt as _; + + let exporter = opentelemetry_sdk::trace::InMemorySpanExporterBuilder::new().build(); + let provider = opentelemetry_sdk::trace::SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let subscriber = tracing_subscriber::registry().with( + tracing_opentelemetry::layer().with_tracer(provider.tracer("vm-driver-test")), + ); + tracing::subscriber::set_global_default(subscriber) + .expect("vm driver test subscriber installs once"); + std::mem::forget(provider); + exporter + }); + static TRACE_LOCK: std::sync::LazyLock> = std::sync::LazyLock::new(|| Mutex::new(())); + + fn assert_is_root(span: &opentelemetry_sdk::trace::SpanData) { + assert_eq!( + span.parent_span_id, + opentelemetry::trace::SpanId::INVALID, + "{:?} should be a trace root", + span.name + ); + } + + fn assert_has_parent(span: &opentelemetry_sdk::trace::SpanData) { + assert_ne!( + span.parent_span_id, + opentelemetry::trace::SpanId::INVALID, + "{:?} should have a parent", + span.name + ); + } + + 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 traced_driver_client( + driver: VmDriver, + ) -> openshell_core::proto::compute::v1::compute_driver_client::ComputeDriverClient< + tonic::transport::Channel, + > { + 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(); + tokio::spawn(async move { + tonic::transport::Server::builder() + .layer(crate::otel_tracing::compute_driver_rpc_layer()) + .add_service(ComputeDriverServer::new(driver)) + .serve_with_incoming(tokio_stream::wrappers::TcpListenerStream::new(listener)) + .await + .unwrap(); + }); + + openshell_core::proto::compute::v1::compute_driver_client::ComputeDriverClient::connect( + format!("http://{address}"), + ) + .await + .unwrap() + } + + #[tokio::test] + async fn compute_driver_rpc_span_continues_the_gateway_trace() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + let mut client = traced_driver_client(driver).await; + + client + .get_capabilities(request_with_traceparent(GetCapabilitiesRequest {})) + .await + .unwrap(); + + let spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + let rpc_spans = spans + .iter() + .filter(|span| span.name == "driver.get_capabilities") + .collect::>(); + assert_eq!( + rpc_spans.len(), + 1, + "middleware should create exactly one VM driver RPC span, got {:?}", + spans.iter().map(|span| &span.name).collect::>() + ); + let span = rpc_spans[0]; + assert_eq!( + span.span_context.trace_id().to_string(), + "4bf92f3577b34da6a3ce929d0e0e4736" + ); + assert_eq!(span.parent_span_id.to_string(), "00f067aa0ba902b7"); + } + + #[tokio::test] + async fn compute_driver_rpcs_record_server_spans_and_error_status() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + let mut client = traced_driver_client(driver).await; + + client + .get_capabilities(request_with_traceparent(GetCapabilitiesRequest {})) + .await + .unwrap(); + assert!( + client + .validate_sandbox_create(request_with_traceparent(ValidateSandboxCreateRequest { + sandbox: None, + })) + .await + .is_err() + ); + assert!( + client + .create_sandbox(request_with_traceparent(CreateSandboxRequest { + sandbox: None, + })) + .await + .is_err() + ); + assert!( + client + .get_sandbox(request_with_traceparent(GetSandboxRequest { + sandbox_id: String::new(), + sandbox_name: String::new(), + })) + .await + .is_err() + ); + client + .list_sandboxes(request_with_traceparent(ListSandboxesRequest {})) + .await + .unwrap(); + assert!( + client + .stop_sandbox(request_with_traceparent(StopSandboxRequest { + sandbox_id: String::new(), + sandbox_name: String::new(), + })) + .await + .is_err() + ); + client + .delete_sandbox(request_with_traceparent(DeleteSandboxRequest { + sandbox_id: String::new(), + sandbox_name: String::new(), + })) + .await + .unwrap(); + let watch = client + .watch_sandboxes(request_with_traceparent(WatchSandboxesRequest {})) + .await + .unwrap(); + drop(watch); + for _ in 0..100 { + if TRACE_EXPORTER + .get_finished_spans() + .unwrap() + .iter() + .any(|span| span.name == "driver.watch_sandboxes") + { + break; + } + tokio::task::yield_now().await; + } + + let spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + let expected = [ + "driver.get_capabilities", + "driver.validate_sandbox_create", + "driver.create_sandbox", + "driver.get_sandbox", + "driver.list_sandboxes", + "driver.stop_sandbox", + "driver.delete_sandbox", + "driver.watch_sandboxes", + ]; + for name in expected { + let span = spans + .iter() + .find(|span| span.name == name) + .unwrap_or_else(|| panic!("missing {name} span")); + assert_eq!(span.span_kind, opentelemetry::trace::SpanKind::Server); + assert_has_parent(span); + } + for name in [ + "driver.validate_sandbox_create", + "driver.create_sandbox", + "driver.get_sandbox", + "driver.stop_sandbox", + ] { + let span = spans.iter().find(|span| span.name == name).unwrap(); + assert!( + matches!(span.status, opentelemetry::trace::Status::Error { .. }), + "{name} should record an error status, got {:?}", + span.status + ); + } + let delete_rpc = spans + .iter() + .find(|span| span.name == "driver.delete_sandbox") + .expect("delete RPC span"); + let cleanup = spans + .iter() + .find(|span| { + span.name == "vm.delete" + && span.span_context.trace_id() == delete_rpc.span_context.trace_id() + }) + .expect("delete cleanup span"); + assert_has_parent(cleanup); + } + + #[tokio::test] + async fn spawned_provisioning_and_phases_have_parents() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let temp = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = temp.path().to_path_buf(); + let sandbox = Sandbox { + id: "sb-spawned-trace".to_string(), + name: "spawned-trace".to_string(), + spec: Some(SandboxSpec { + template: Some(SandboxTemplate { + image: "invalid image reference".to_string(), + ..Default::default() + }), + ..Default::default() + }), + ..Default::default() + }; + let request = request_with_traceparent(CreateSandboxRequest { + sandbox: Some(sandbox), + }); + + let mut client = traced_driver_client(driver).await; + client.create_sandbox(request).await.unwrap(); + + let mut spans = Vec::new(); + for _ in 0..100 { + spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + if spans.iter().any(|span| span.name == "vm.provision") { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + let provisioning = spans + .iter() + .find(|span| span.name == "vm.provision") + .expect("spawned provisioning span"); + assert_has_parent(provisioning); + let prepare_images = spans + .iter() + .find(|span| span.name == "vm.prepare_images") + .expect("image preparation span"); + assert_has_parent(prepare_images); + let resolve_bootstrap = spans + .iter() + .find(|span| span.name == "vm.resolve_bootstrap_image") + .expect("bootstrap image resolution span"); + assert_has_parent(resolve_bootstrap); + } + + #[tokio::test] + async fn startup_reconciliation_is_root_and_restore_operations_have_parents() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let temp = tempfile::tempdir().unwrap(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.state_dir = temp.path().to_path_buf(); + for suffix in ["a", "b"] { + let sandbox = Sandbox { + id: format!("sb-restored-trace-{suffix}"), + name: format!("restored-trace-{suffix}"), + spec: Some(SandboxSpec { + template: Some(SandboxTemplate { + image: "invalid image reference".to_string(), + ..Default::default() + }), + ..Default::default() + }), + ..Default::default() + }; + let state_dir = temp.path().join("sandboxes").join(&sandbox.id); + tokio::fs::create_dir_all(&state_dir).await.unwrap(); + write_sandbox_request(&state_dir, &sandbox).await.unwrap(); + } + + driver.restore_persisted_sandboxes().await; + + let mut spans = Vec::new(); + for _ in 0..100 { + spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + if spans.iter().any(|span| span.name == "reconcile.sandboxes") { + break; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + + let reconciliations = spans + .iter() + .filter(|span| span.name == "reconcile.sandboxes") + .collect::>(); + assert_eq!( + reconciliations.len(), + 1, + "startup should reconcile all persisted sandboxes in one trace" + ); + let reconciliation = reconciliations[0]; + assert_is_root(reconciliation); + let restorations = spans + .iter() + .filter(|span| span.name == "vm.restore") + .collect::>(); + assert_eq!(restorations.len(), 2); + for restoration in restorations { + assert_has_parent(restoration); + } + let provisioning = spans + .iter() + .filter(|span| span.name == "vm.provision") + .collect::>(); + assert_eq!(provisioning.len(), 2); + for span in provisioning { + assert_has_parent(span); + } + let prepare_images = spans + .iter() + .filter(|span| span.name == "vm.prepare_images") + .collect::>(); + assert_eq!(prepare_images.len(), 2); + for span in prepare_images { + assert_has_parent(span); + } + } + + #[tokio::test] + async fn background_provisioning_does_not_extend_the_rpc_span_lifetime() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let rpc = tracing::info_span!("driver.create_sandbox"); + let entered = rpc.enter(); + let provisioning = + provisioning_span(&rpc.context(), "sb-lifetime", "invalid image reference"); + drop(entered); + drop(rpc); + + assert!( + TRACE_EXPORTER + .get_finished_spans() + .unwrap() + .iter() + .any(|span| span.name == "driver.create_sandbox"), + "the RPC span should finish while background provisioning is still active" + ); + drop(provisioning); + } + + #[tokio::test] + async fn overlay_preparation_records_a_provisioning_phase_span() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let mut driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + driver.config.overlay_disk_mib = u64::MAX; + let parent = tracing::info_span!("vm.provision"); + + let result = driver + .prepare_runtime_overlay(Path::new("/unused"), None, None, OverlayPreparation::Fresh) + .instrument(parent) + .await; + assert!(result.is_err(), "overflow should stop before disk I/O"); + + let spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + let overlay = spans + .iter() + .find(|span| span.name == "vm.prepare_overlay") + .expect("overlay preparation span"); + assert_has_parent(overlay); + } + + #[tokio::test] + async fn post_overlay_provisioning_stages_record_child_spans() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let driver = test_driver_with_extensions(LifecycleExtensionRegistry::new()); + let sandbox = Sandbox { + id: "sb-post-overlay".to_string(), + ..Default::default() + }; + let mut plan = driver + .build_vm_launch_plan(&sandbox.id, false, false, None) + .unwrap(); + let provisioning = tracing::info_span!("vm.provision"); + + async { + driver + .lifecycle_extensions + .configure_launch(&sandbox, Path::new("/unused"), &mut plan) + .await + .unwrap(); + driver + .lifecycle_extensions + .before_launch(&sandbox, Path::new("/unused"), &mut plan) + .await + .unwrap(); + let invalid_dropin = GuestInitDropin::new("../invalid", Vec::new()); + assert!( + inject_guest_init_dropins(Path::new("/unused"), &[invalid_dropin]).is_err(), + "an invalid drop-in should fail after creating its span" + ); + } + .instrument(provisioning) + .await; + + let spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + for name in [ + "vm.configure_launch", + "vm.before_launch", + "vm.prepare_guest", + ] { + let span = spans + .iter() + .find(|span| span.name == name) + .unwrap_or_else(|| panic!("missing {name} span")); + assert_has_parent(span); + } + } + + #[tokio::test] + async fn launcher_spawn_records_a_provisioning_phase_span() { + let _lock = TRACE_LOCK.lock().await; + TRACE_EXPORTER.reset(); + let provisioning = tracing::info_span!("vm.provision"); + let mut command = Command::new("sh"); + command.arg("-c").arg("exit 0"); + + let mut child = async { + spawn_vm_launcher(&mut command, "sb-launch-trace", &VmBackend::Libkrun).unwrap() + } + .instrument(provisioning) + .await; + child.wait().await.unwrap(); + + let spans = TRACE_EXPORTER.get_finished_spans().unwrap(); + let launch = spans + .iter() + .find(|span| span.name == "vm.launch") + .expect("launcher span"); + assert_has_parent(launch); + } fn gpu_device_ids_config(device_ids: &[&str]) -> Struct { list_string_driver_config("gpu_device_ids", device_ids) diff --git a/crates/openshell-driver-vm/src/lib.rs b/crates/openshell-driver-vm/src/lib.rs index 88e2c3b201..98ba6b0c9a 100644 --- a/crates/openshell-driver-vm/src/lib.rs +++ b/crates/openshell-driver-vm/src/lib.rs @@ -7,6 +7,7 @@ mod ffi; pub mod gpu; pub mod lifecycle; mod nft_ruleset; +pub mod otel_tracing; pub mod procguard; mod rootfs; mod runtime; diff --git a/crates/openshell-driver-vm/src/lifecycle.rs b/crates/openshell-driver-vm/src/lifecycle.rs index c042715e58..230eb4071a 100644 --- a/crates/openshell-driver-vm/src/lifecycle.rs +++ b/crates/openshell-driver-vm/src/lifecycle.rs @@ -550,6 +550,15 @@ impl LifecycleExtensionRegistry { .collect() } + #[tracing::instrument( + name = "vm.configure_launch", + skip_all, + fields( + otel.name = "vm.configure_launch", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox.id, + ) + )] pub async fn configure_launch( &self, sandbox: &Sandbox, @@ -588,6 +597,15 @@ impl LifecycleExtensionRegistry { Ok(()) } + #[tracing::instrument( + name = "vm.before_launch", + skip_all, + fields( + otel.name = "vm.before_launch", + otel.status_code = tracing::field::Empty, + sandbox.id = %sandbox.id, + ) + )] pub async fn before_launch( &self, sandbox: &Sandbox, diff --git a/crates/openshell-driver-vm/src/main.rs b/crates/openshell-driver-vm/src/main.rs index 0ae694effa..40d37e120b 100644 --- a/crates/openshell-driver-vm/src/main.rs +++ b/crates/openshell-driver-vm/src/main.rs @@ -6,6 +6,7 @@ use futures::Stream; use miette::{IntoDiagnostic, Result}; use openshell_core::VERSION; use openshell_core::proto::compute::v1::compute_driver_server::ComputeDriverServer; +use openshell_driver_vm::otel_tracing::compute_driver_rpc_layer; #[cfg(target_os = "macos")] use openshell_driver_vm::{VM_RUNTIME_DIR_ENV, configured_runtime_dir}; use openshell_driver_vm::{VmBackend, VmDriver, VmDriverConfig, VmLaunchConfig, procguard, run_vm}; @@ -18,6 +19,7 @@ use std::task::{Context, Poll}; use tokio::net::{UnixListener, UnixStream}; use tracing::info; use tracing_subscriber::EnvFilter; +use tracing_subscriber::prelude::*; #[derive(Parser, Debug)] #[command(name = "openshell-driver-vm")] @@ -86,6 +88,9 @@ struct Args { #[arg(long, env = "OPENSHELL_LOG_LEVEL", default_value = "info")] log_level: String, + #[arg(long, env = "OPENSHELL_OTLP_ENDPOINT")] + otlp_endpoint: Option, + #[arg(long, env = "OPENSHELL_GRPC_ENDPOINT")] openshell_endpoint: Option, @@ -181,11 +186,22 @@ async fn main() -> Result<()> { return Ok(()); } - 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_vm::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_vm::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 listen_mode = compute_driver_listen_mode(&args).map_err(|err| miette::miette!("{err}"))?; @@ -226,7 +242,7 @@ async fn main() -> Result<()> { .await .map_err(|err| miette::miette!("{err}"))?; - match listen_mode { + let result = match listen_mode { ComputeDriverListenMode::Unix { socket_path, expected_peer_pid, @@ -237,6 +253,7 @@ async fn main() -> Result<()> { let listener = UnixListener::bind(&socket_path).into_diagnostic()?; restrict_socket_permissions(&socket_path).map_err(|err| miette::miette!("{err}"))?; let result = tonic::transport::Server::builder() + .layer(compute_driver_rpc_layer()) .add_service(ComputeDriverServer::new(driver)) .serve_with_incoming(AuthenticatedUnixIncoming::new(listener, expected_peer_pid)) .await @@ -247,12 +264,19 @@ async fn main() -> Result<()> { ComputeDriverListenMode::Tcp(bind_address) => { info!(address = %bind_address, "Starting unauthenticated dev vm compute driver"); tonic::transport::Server::builder() + .layer(compute_driver_rpc_layer()) .add_service(ComputeDriverServer::new(driver)) .serve(bind_address) .await .into_diagnostic() } + }; + if let Some(provider) = &tracer_provider + && let Err(error) = provider.shutdown() + { + tracing::warn!(%error, "OTLP tracer provider shutdown failed"); } + result } #[derive(Debug, Clone, PartialEq, Eq)] @@ -649,6 +673,19 @@ mod tests { assert!(err.contains("--bind-socket is required")); } + #[test] + fn accepts_gateway_otlp_endpoint() { + let args = Args::try_parse_from([ + "openshell-driver-vm", + "--otlp-endpoint", + "http://127.0.0.1:4317", + ]); + assert!( + args.is_ok(), + "VM driver should accept the gateway OTLP endpoint" + ); + } + #[test] fn listen_mode_rejects_bind_address_without_tcp_opt_in() { let args = Args::parse_from(["openshell-driver-vm", "--bind-address", "127.0.0.1:50061"]); diff --git a/crates/openshell-driver-vm/src/otel_tracing.rs b/crates/openshell-driver-vm/src/otel_tracing.rs new file mode 100644 index 0000000000..d71391e6ea --- /dev/null +++ b/crates/openshell-driver-vm/src/otel_tracing.rs @@ -0,0 +1,242 @@ +// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +//! OpenTelemetry trace exporting. + +use http::{HeaderMap, Request}; +use openshell_otel::{OtlpTraceConfig, SdkTracerProvider, ServiceName, SetupError}; +use opentelemetry::propagation::{Extractor, TextMapPropagator}; +use opentelemetry::trace::TraceContextExt as _; +use opentelemetry_sdk::propagation::TraceContextPropagator; +use tower_http::classify::GrpcFailureClass; +use tower_http::trace::{GrpcMakeClassifier, MakeSpan, OnFailure, TraceLayer}; +use tracing::{Span, Subscriber}; +use tracing_opentelemetry::OpenTelemetrySpanExt as _; +use tracing_subscriber::registry::LookupSpan; + +const SERVICE_NAME: &str = "openshell-driver-vm"; +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 { + TraceLayer::new_for_grpc() + .make_span_with(ComputeDriverRpcSpan) + .on_request(()) + .on_response(()) + .on_body_chunk(()) + .on_eos(()) + .on_failure(RecordGrpcFailure) +} + +/// Creates the server span for an inbound compute-driver request. +#[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(), + &HeaderExtractor(request.headers()), + ); + if parent.span().span_context().is_valid() { + let _ = span.set_parent(parent); + } + span + } +} + +fn compute_driver_rpc_operation(path: &str) -> (&'static str, &'static str) { + match path.rsplit('/').next() { + Some("GetCapabilities") => ("driver.get_capabilities", "get_capabilities"), + 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("DeleteSandbox") => ("driver.delete_sandbox", "delete_sandbox"), + Some("WatchSandboxes") => ("driver.watch_sandboxes", "watch_sandboxes"), + _ => ("driver.unknown", "unknown"), + } +} + +struct HeaderExtractor<'a>(&'a HeaderMap); + +impl Extractor for HeaderExtractor<'_> { + fn get(&self, key: &str) -> Option<&str> { + self.0.get(key).and_then(|value| value.to_str().ok()) + } + + fn keys(&self) -> Vec<&str> { + self.0.keys().map(http::HeaderName::as_str).collect() + } +} + +/// Records a non-OK gRPC outcome on the request span. +#[derive(Debug, Clone, Copy)] +pub struct RecordGrpcFailure; + +impl OnFailure for RecordGrpcFailure { + fn on_failure( + &mut self, + failure: GrpcFailureClass, + _latency: std::time::Duration, + span: &Span, + ) { + span.record("otel.status_code", "ERROR"); + if let GrpcFailureClass::Code(code) = failure { + span.record("rpc.grpc.status_code", code.get()); + } + } +} + +/// Build a tracer provider for the configured OTLP/gRPC endpoint. +#[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(), + })) +} + +/// Build the tracing layer that exports VM-driver spans. +pub fn layer(provider: &SdkTracerProvider) -> openshell_otel::OtlpLayer +where + S: Subscriber + for<'span> LookupSpan<'span>, +{ + openshell_otel::layer(provider, INSTRUMENTATION_SCOPE) +} + +#[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>, + } + + #[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); + } + } + Ok(tonic::Response::new(ExportTraceServiceResponse::default())) + } + } + + #[test] + fn compute_driver_rpc_names_are_explicitly_mapped_and_schema_bounded() { + let list_sandboxes: (&'static str, &'static str) = super::compute_driver_rpc_operation( + "/openshell.compute.v1.ComputeDriver/ListSandboxes", + ); + assert_eq!(list_sandboxes, ("driver.list_sandboxes", "list_sandboxes")); + assert_eq!( + super::compute_driver_rpc_operation( + "/openshell.compute.v1.ComputeDriver/AttackerControlled12345" + ), + ("driver.unknown", "unknown"), + "paths absent from the protobuf schema must not create span names" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] + async fn vm_driver_spans_reach_otlp_collector_with_distinct_service_name() { + let received = Arc::new(Mutex::new(Received::default())); + 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), + }; + tokio::spawn(async move { + tonic::transport::Server::builder() + .add_service(TraceServiceServer::new(collector)) + .serve_with_incoming(tokio_stream::wrappers::TcpListenerStream::new(listener)) + .await + .unwrap(); + }); + + let (provider, error) = super::provider_for(Some(&format!("http://{address}"))); + assert!(error.is_none(), "valid OTLP endpoint should configure"); + let provider = provider.expect("provider"); + let subscriber = tracing_subscriber::registry().with(super::layer(&provider)); + tracing::subscriber::with_default(subscriber, || { + let span = tracing::info_span!("vm.provision", sandbox.id = "sb-otlp"); + drop(span.enter()); + drop(span); + }); + provider.force_flush().unwrap(); + + for _ in 0..100 { + if !received.lock().unwrap().spans.is_empty() { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + let received = received.lock().unwrap(); + received + .spans + .iter() + .find(|span| span.name == "vm.provision") + .expect("VM span should reach collector"); + assert!( + received + .service_names + .iter() + .any(|name| name == "openshell-driver-vm"), + "VM spans should use a distinct service name, got {:?}", + received.service_names + ); + provider.shutdown().ok(); + } +} diff --git a/crates/openshell-server/src/compute/mod.rs b/crates/openshell-server/src/compute/mod.rs index 25a2655a74..4a37580daa 100644 --- a/crates/openshell-server/src/compute/mod.rs +++ b/crates/openshell-server/src/compute/mod.rs @@ -13,6 +13,7 @@ pub use openshell_driver_podman::PodmanComputeConfig; pub use vm::VmComputeConfig; use crate::grpc::policy::SANDBOX_SETTINGS_OBJECT_TYPE; +use crate::otel_tracing::TraceContextInterceptor; use crate::persistence::{ DRAFT_CHUNK_OBJECT_TYPE, ObjectId, ObjectName, ObjectRecord, ObjectType, POLICY_OBJECT_TYPE, Store, WriteCondition, @@ -364,16 +365,22 @@ impl AcquiredRemoteDriverEndpoint { #[derive(Debug, Clone)] struct RemoteComputeDriver { - channel: Channel, + client: RemoteComputeDriverClient, } +type RemoteComputeDriverClient = ComputeDriverClient< + tonic::service::interceptor::InterceptedService, +>; + impl RemoteComputeDriver { fn new(channel: Channel) -> Self { - Self { channel } + Self { + client: ComputeDriverClient::with_interceptor(channel, TraceContextInterceptor), + } } - fn client(&self) -> ComputeDriverClient { - ComputeDriverClient::new(self.channel.clone()) + fn client(&self) -> RemoteComputeDriverClient { + self.client.clone() } } @@ -6667,6 +6674,134 @@ mod tests { test_exporter::assert_is_root(&initialization); } + #[tokio::test] + #[cfg(unix)] + async fn remote_compute_driver_interceptor_propagates_every_rpc() { + use crate::otel_tracing::test_collector; + use crate::test_support::FakeComputeDriver; + + let dir = tempfile::tempdir().unwrap(); + let socket_path = dir.path().join("compute-driver.sock"); + let driver = FakeComputeDriver::new(); + let _server = driver.serve_uds(&socket_path).unwrap(); + let endpoint = connect_remote_compute_driver("external-test", &socket_path) + .await + .unwrap(); + let remote = RemoteComputeDriver::new(endpoint.channel); + let sandbox = DriverSandbox { + id: "sb-trace".to_string(), + name: "trace-sandbox".to_string(), + ..Default::default() + }; + + let traced = test_collector::install_traced(); + async { + remote + .get_capabilities(Request::new(GetCapabilitiesRequest {})) + .await + .unwrap(); + remote + .validate_sandbox_create(Request::new(ValidateSandboxCreateRequest { + sandbox: Some(sandbox.clone()), + })) + .await + .unwrap(); + remote + .create_sandbox(Request::new(CreateSandboxRequest { + sandbox: Some(sandbox.clone()), + })) + .await + .unwrap(); + remote + .get_sandbox(Request::new(GetSandboxRequest { + sandbox_id: sandbox.id.clone(), + sandbox_name: String::new(), + })) + .await + .unwrap(); + remote + .list_sandboxes(Request::new(ListSandboxesRequest {})) + .await + .unwrap(); + remote + .stop_sandbox(Request::new(StopSandboxRequest { + sandbox_id: sandbox.id.clone(), + sandbox_name: String::new(), + })) + .await + .unwrap(); + remote + .watch_sandboxes(Request::new(WatchSandboxesRequest {})) + .await + .unwrap(); + remote + .delete_sandbox(Request::new(DeleteSandboxRequest { + sandbox_id: sandbox.id, + sandbox_name: String::new(), + })) + .await + .unwrap(); + } + .instrument(tracing::info_span!("request")) + .await; + + let request_spans = traced.spans_named("request"); + assert_eq!(request_spans.len(), 1, "one request span should finish"); + let trace_id = request_spans[0].span_context.trace_id().to_string(); + let traceparents = driver.traceparents(); + assert_eq!( + traceparents.len(), + 8, + "the client interceptor should cover every RPC" + ); + assert!( + traceparents + .iter() + .all(|traceparent| traceparent.contains(&trace_id)), + "every RPC should carry the active trace ID; got {traceparents:?}" + ); + } + + #[tokio::test] + #[cfg(unix)] + async fn remote_compute_driver_initialization_parents_the_capability_probe() { + use crate::otel_tracing::test_collector; + use crate::test_support::FakeComputeDriver; + + let dir = tempfile::tempdir().unwrap(); + let socket_path = dir.path().join("compute-driver.sock"); + let driver = FakeComputeDriver::new(); + let _server = driver.serve_uds(&socket_path).unwrap(); + let endpoint = connect_remote_compute_driver("external-test", &socket_path) + .await + .unwrap(); + let store = Arc::new(Store::connect("sqlite::memory:").await.unwrap()); + + let traced = test_collector::install_traced(); + ComputeRuntime::new_remote_driver( + endpoint, + store, + SandboxIndex::new(), + SandboxWatchBus::new(), + TracingLogBus::new(), + Arc::new(SupervisorSessionRegistry::new()), + ) + .await + .unwrap(); + + let initialization = traced.span_with("driver.initialize", "driver.name", "external-test"); + let trace_id = initialization.span_context.trace_id().to_string(); + assert_eq!( + driver.traceparents().len(), + 1, + "the capability probe should carry initialization trace context" + ); + assert!( + driver.traceparents()[0].contains(&trace_id), + "the capability probe should be part of the initialization trace" + ); + } + #[tokio::test] #[cfg(unix)] async fn remote_compute_driver_forwards_lifecycle_calls_over_uds() { diff --git a/crates/openshell-server/src/compute/vm.rs b/crates/openshell-server/src/compute/vm.rs index be88047f33..918309fc4c 100644 --- a/crates/openshell-server/src/compute/vm.rs +++ b/crates/openshell-server/src/compute/vm.rs @@ -32,6 +32,9 @@ use super::AcquiredRemoteDriverEndpoint; #[cfg(unix)] use super::ManagedDriverProcess; +use crate::config_file::OtlpConfig; +#[cfg(unix)] +use crate::otel_tracing::TraceContextInterceptor; #[cfg(unix)] use hyper_util::rt::TokioIo; #[cfg(unix)] @@ -452,6 +455,7 @@ pub fn compute_driver_guest_tls_paths( pub async fn spawn( config: &Config, vm_config: &VmComputeConfig, + otlp_config: Option<&OtlpConfig>, ) -> Result { if vm_config.grpc_endpoint.trim().is_empty() { return Err(Error::config( @@ -474,6 +478,7 @@ pub async fn spawn( .arg("--expected-peer-pid") .arg(std::process::id().to_string()); command.arg("--log-level").arg(&config.log_level); + append_otlp_args(&mut command, otlp_config); command .arg("--openshell-endpoint") .arg(&vm_config.grpc_endpoint); @@ -515,10 +520,18 @@ pub async fn spawn( )) } +#[cfg(unix)] +fn append_otlp_args(command: &mut Command, otlp_config: Option<&OtlpConfig>) { + if let Some(config) = otlp_config { + command.arg("--otlp-endpoint").arg(&config.endpoint); + } +} + #[cfg(not(unix))] pub async fn spawn( _config: &Config, _vm_config: &VmComputeConfig, + _otlp_config: Option<&OtlpConfig>, ) -> Result { Err(Error::config( "the vm compute driver requires unix domain socket support", @@ -526,6 +539,15 @@ pub async fn spawn( } #[cfg(unix)] +#[tracing::instrument( + name = "driver.wait_for_ready", + skip_all, + fields( + otel.name = "driver.wait_for_ready", + otel.status_code = tracing::field::Empty, + driver.name = "vm", + ) +)] async fn wait_for_compute_driver( socket_path: &Path, child: &mut tokio::process::Child, @@ -543,7 +565,8 @@ async fn wait_for_compute_driver( match connect_compute_driver(socket_path).await { Ok(channel) => { - let mut client = ComputeDriverClient::new(channel.clone()); + let mut client = + ComputeDriverClient::with_interceptor(channel.clone(), TraceContextInterceptor); match client .get_capabilities(tonic::Request::new(GetCapabilitiesRequest {})) .await @@ -586,15 +609,72 @@ async fn connect_compute_driver(socket_path: &Path) -> Result { #[cfg(all(test, unix))] mod tests { use super::{ - VmComputeConfig, compute_driver_guest_tls_paths, compute_driver_socket_path, current_euid, - prepare_compute_driver_socket_path, prepare_vm_state_dir, resolve_compute_driver_bin, - resolve_driver_search_dirs, + VmComputeConfig, append_otlp_args, compute_driver_guest_tls_paths, + compute_driver_socket_path, current_euid, prepare_compute_driver_socket_path, + prepare_vm_state_dir, resolve_compute_driver_bin, resolve_driver_search_dirs, + wait_for_compute_driver, }; + use crate::config_file::OtlpConfig; use std::os::unix::fs::PermissionsExt; use std::os::unix::net::UnixListener as StdUnixListener; use std::path::PathBuf; use tempfile::tempdir; + #[test] + fn vm_driver_command_includes_gateway_otlp_endpoint() { + let mut command = tokio::process::Command::new("openshell-driver-vm"); + append_otlp_args( + &mut command, + Some(&OtlpConfig { + endpoint: "http://collector.internal:4317".to_string(), + service_name: Some("custom-gateway".to_string()), + }), + ); + + let args = command + .as_std() + .get_args() + .map(|arg| arg.to_string_lossy().into_owned()) + .collect::>(); + assert_eq!(args, ["--otlp-endpoint", "http://collector.internal:4317"]); + } + + #[tokio::test] + async fn readiness_probe_propagates_the_active_trace() { + use crate::otel_tracing::test_collector; + use crate::test_support::FakeComputeDriver; + + let dir = tempdir().unwrap(); + let socket_path = dir.path().join("compute-driver.sock"); + let driver = FakeComputeDriver::new(); + let _server = driver.serve_uds(&socket_path).unwrap(); + let mut child = tokio::process::Command::new("sh") + .arg("-c") + .arg("sleep 30") + .kill_on_drop(true) + .spawn() + .unwrap(); + + let traced = test_collector::install_traced(); + wait_for_compute_driver(&socket_path, &mut child) + .await + .unwrap(); + + let readiness = traced.spans_named("driver.wait_for_ready"); + assert_eq!(readiness.len(), 1, "one readiness operation should finish"); + test_collector::assert_is_root(&readiness[0]); + let trace_id = readiness[0].span_context.trace_id().to_string(); + assert_eq!( + driver.traceparents().len(), + 1, + "the readiness capability probe should carry trace context" + ); + assert!( + driver.traceparents()[0].contains(&trace_id), + "the readiness probe should be part of the active trace" + ); + } + #[test] fn resolve_driver_bin_uses_driver_dir_when_binary_present() { let dir = tempdir().unwrap(); diff --git a/crates/openshell-server/src/lib.rs b/crates/openshell-server/src/lib.rs index 1ab9e1ada1..6151922b71 100644 --- a/crates/openshell-server/src/lib.rs +++ b/crates/openshell-server/src/lib.rs @@ -849,7 +849,10 @@ async fn build_compute_runtime( } ConfiguredComputeDriver::Builtin(ComputeDriverKind::Vm) => { let vm_config = compute::driver_config::vm_config_from_context(driver_startup)?; - let endpoint = compute::vm::spawn(config, &vm_config).await?; + let otlp_config = driver_startup + .file + .and_then(|file| file.openshell.gateway.otlp.as_ref()); + let endpoint = compute::vm::spawn(config, &vm_config, otlp_config).await?; ComputeRuntime::new_remote_driver( endpoint, store, diff --git a/crates/openshell-server/src/otel_tracing.rs b/crates/openshell-server/src/otel_tracing.rs index 58b4cdf804..d050999c21 100644 --- a/crates/openshell-server/src/otel_tracing.rs +++ b/crates/openshell-server/src/otel_tracing.rs @@ -23,10 +23,13 @@ pub use openshell_otel::SetupError; use openshell_otel::{OtlpTraceConfig, ServiceName}; +use opentelemetry::propagation::{Injector, TextMapPropagator}; #[cfg(test)] use opentelemetry_sdk::Resource; +use opentelemetry_sdk::propagation::TraceContextPropagator; use opentelemetry_sdk::trace::SdkTracerProvider; use tracing::Subscriber; +use tracing_opentelemetry::OpenTelemetrySpanExt as _; use tracing_subscriber::registry::LookupSpan; use crate::config_file::OtlpConfig; @@ -37,6 +40,36 @@ const DEFAULT_SERVICE_NAME: &str = "openshell-gateway"; /// Instrumentation scope recorded on spans this gateway emits. const INSTRUMENTATION_SCOPE: &str = "openshell-gateway"; +/// Inject the active trace context into an outbound tonic request. +#[derive(Debug, Clone, Copy)] +pub struct TraceContextInterceptor; + +impl tonic::service::Interceptor for TraceContextInterceptor { + fn call( + &mut self, + mut request: tonic::Request<()>, + ) -> Result, tonic::Status> { + let context = tracing::Span::current().context(); + TraceContextPropagator::new() + .inject_context(&context, &mut MetadataInjector(request.metadata_mut())); + Ok(request) + } +} + +struct MetadataInjector<'a>(&'a mut tonic::metadata::MetadataMap); + +impl Injector for MetadataInjector<'_> { + fn set(&mut self, key: &str, value: String) { + let Ok(key) = key.parse::>() else { + return; + }; + let Ok(value) = value.parse() else { + return; + }; + self.0.insert(key, value); + } +} + fn trace_config(cfg: &OtlpConfig) -> OtlpTraceConfig<'_> { let service_name = cfg .service_name diff --git a/crates/openshell-server/src/test_support.rs b/crates/openshell-server/src/test_support.rs index 2bfa9998a2..c3327d428b 100644 --- a/crates/openshell-server/src/test_support.rs +++ b/crates/openshell-server/src/test_support.rs @@ -72,6 +72,7 @@ struct FakeComputeDriverState { gateway_listener_requirements_supported: bool, sandboxes: HashMap, calls: Vec, + traceparents: Vec, } impl Default for FakeComputeDriver { @@ -92,6 +93,7 @@ impl FakeComputeDriver { gateway_listener_requirements_supported: true, sandboxes: HashMap::new(), calls: Vec::new(), + traceparents: Vec::new(), })), } } @@ -142,6 +144,11 @@ impl FakeComputeDriver { self.with_state(|state| state.calls.clone()) } + #[must_use] + pub fn traceparents(&self) -> Vec { + self.with_state(|state| state.traceparents.clone()) + } + pub fn clear_calls(&self) { self.with_state(|state| state.calls.clear()); } @@ -170,6 +177,14 @@ impl FakeComputeDriver { .expect("fake compute driver state poisoned"); f(&mut state) } + + fn record_traceparent(&self, metadata: &tonic::metadata::MetadataMap) { + let traceparent = metadata + .get("traceparent") + .and_then(|value| value.to_str().ok()) + .map(str::to_string); + self.with_state(|state| state.traceparents.extend(traceparent)); + } } #[cfg(unix)] @@ -211,8 +226,9 @@ impl ComputeDriver for FakeComputeDriver { async fn get_capabilities( &self, - _request: Request, + request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let response = self.with_state(|state| { state.calls.push(FakeComputeDriverCall::GetCapabilities); GetCapabilitiesResponse { @@ -246,6 +262,7 @@ impl ComputeDriver for FakeComputeDriver { &self, request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let sandbox = request.into_inner().sandbox; self.with_state(|state| { state @@ -259,6 +276,7 @@ impl ComputeDriver for FakeComputeDriver { &self, request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let request = request.into_inner(); let sandbox = self.with_state(|state| { state.calls.push(FakeComputeDriverCall::GetSandbox { @@ -283,8 +301,9 @@ impl ComputeDriver for FakeComputeDriver { async fn list_sandboxes( &self, - _request: Request, + request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let sandboxes = self.with_state(|state| { state.calls.push(FakeComputeDriverCall::ListSandboxes); state.sandboxes.values().cloned().collect() @@ -296,6 +315,7 @@ impl ComputeDriver for FakeComputeDriver { &self, request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let sandbox = request.into_inner().sandbox; self.with_state(|state| { if let Some(sandbox) = sandbox.as_ref() { @@ -312,6 +332,7 @@ impl ComputeDriver for FakeComputeDriver { &self, request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let request = request.into_inner(); self.with_state(|state| { state.calls.push(FakeComputeDriverCall::StopSandbox { @@ -326,6 +347,7 @@ impl ComputeDriver for FakeComputeDriver { &self, request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); let request = request.into_inner(); let deleted = self.with_state(|state| { state.calls.push(FakeComputeDriverCall::DeleteSandbox { @@ -351,8 +373,9 @@ impl ComputeDriver for FakeComputeDriver { async fn watch_sandboxes( &self, - _request: Request, + request: Request, ) -> Result, Status> { + self.record_traceparent(request.metadata()); self.with_state(|state| state.calls.push(FakeComputeDriverCall::WatchSandboxes)); Ok(Response::new(Box::pin(stream::empty()))) } diff --git a/docs/reference/gateway-config.mdx b/docs/reference/gateway-config.mdx index e46bc7c41c..03f24929d3 100644 --- a/docs/reference/gateway-config.mdx +++ b/docs/reference/gateway-config.mdx @@ -202,10 +202,12 @@ The transport is **OTLP over gRPC only**. HTTP/protobuf and HTTP/JSON are not su The OpenTelemetry SDK logs export failures after startup. Spans in a failed batch are dropped rather than retried. -`service_name` sets the `service.name` resource attribute and defaults to `openshell-gateway`. The gateway also reports `service.version`. +`service_name` sets the gateway's `service.name` resource attribute and defaults to `openshell-gateway`. The gateway also reports `service.version`. Only OpenTelemetry traces are exported. Inbound gRPC and HTTP requests produce server spans named for the RPC or HTTP method. Store and compute-driver operations appear as child spans. Internal reconciliation, credential-refresh, and driver-watch loops create operation roots for their store work because no inbound request supplies a parent. The gateway continues valid W3C `traceparent` context and starts a new trace when none is supplied. Request spans carry `method`, `path`, and the `request_id` that also appears in gateway logs. Health endpoint spans use DEBUG level and are not exported by the default INFO filter. +The gateway forwards the OTLP configuration to managed external drivers. Each driver exports under its own service name. + ### Tuning This table decides whether and where to export. How the SDK exports is controlled by the standard OpenTelemetry environment variables, which the gateway reads through the SDK rather than mirroring as TOML keys: