diff --git a/rust/Cargo.lock b/rust/Cargo.lock index 1493357..fd495df 100644 --- a/rust/Cargo.lock +++ b/rust/Cargo.lock @@ -952,6 +952,15 @@ dependencies = [ "tempfile", ] +[[package]] +name = "nu-ansi-term" +version = "0.50.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" +dependencies = [ + "windows-sys 0.60.2", +] + [[package]] name = "num-traits" version = "0.2.19" @@ -1043,6 +1052,18 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry-appender-tracing" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" +dependencies = [ + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", +] + [[package]] name = "opentelemetry-http" version = "0.32.0" @@ -1749,6 +1770,15 @@ dependencies = [ "serde", ] +[[package]] +name = "sharded-slab" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f40ca3c46823713e0d4209592e8d6e826aa57e928f09752619fc696c499637f6" +dependencies = [ + "lazy_static", +] + [[package]] name = "shlex" version = "1.3.0" @@ -1818,6 +1848,7 @@ dependencies = [ "http", "once_cell", "opentelemetry", + "opentelemetry-appender-tracing", "opentelemetry-http", "opentelemetry-otlp", "opentelemetry-semantic-conventions", @@ -1833,6 +1864,8 @@ dependencies = [ "tokio", "tower-layer", "tower-service", + "tracing", + "tracing-subscriber", "uuid", "wiremock", ] @@ -1944,6 +1977,15 @@ dependencies = [ "syn", ] +[[package]] +name = "thread_local" +version = "1.1.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ad99c4c6d32803332c548b1af0540b357b3f5fc0be8f6c6bfe8b2e6ae784070" +dependencies = [ + "cfg-if", +] + [[package]] name = "tinystr" version = "0.8.3" @@ -2115,6 +2157,32 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "db97caf9d906fbde555dd62fa95ddba9eecfd14cb388e4f491a66d74cd5fb79a" dependencies = [ "once_cell", + "valuable", +] + +[[package]] +name = "tracing-log" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee855f1f400bd0e5c02d150ae5de3840039a3f54b025156404e34c23c03f47c3" +dependencies = [ + "log", + "once_cell", + "tracing-core", +] + +[[package]] +name = "tracing-subscriber" +version = "0.3.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb7f578e5945fb242538965c2d0b04418d38ec25c79d160cd279bf0731c8d319" +dependencies = [ + "nu-ansi-term", + "sharded-slab", + "smallvec", + "thread_local", + "tracing-core", + "tracing-log", ] [[package]] @@ -2176,6 +2244,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "valuable" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" + [[package]] name = "vcpkg" version = "0.2.15" diff --git a/rust/observability/Cargo.toml b/rust/observability/Cargo.toml index d096a48..b8cf23b 100644 --- a/rust/observability/Cargo.toml +++ b/rust/observability/Cargo.toml @@ -47,7 +47,7 @@ thiserror = "2" # opentelemetry-http is pulled with no features — we only need the `HttpClient` # trait to implement it ourselves (the `reqwest` feature, which impls it for # reqwest::Client, is no longer used). -opentelemetry = { version = "0.32", features = ["trace", "metrics"] } +opentelemetry = { version = "0.32", features = ["trace", "metrics", "logs"] } # The two `experimental_*_with_async_runtime` features expose the async-runtime # batch span processor + periodic metric reader used in `otel.rs`. They are # REQUIRED for correctness, not a nicety: the default (thread-based) processors @@ -59,11 +59,16 @@ opentelemetry = { version = "0.32", features = ["trace", "metrics"] } opentelemetry_sdk = { version = "0.32", features = [ "trace", "metrics", + "logs", "rt-tokio", "experimental_trace_batch_span_processor_with_async_runtime", "experimental_metrics_periodicreader_with_async_runtime", + # Same reactor-panic rationale as traces/metrics (SMOODEV-2045): the DEFAULT + # `BatchLogProcessor` drives export with `block_on` on a bare OS thread where + # our reqwest-backed exporter has no Tokio reactor and aborts the process. + "experimental_logs_batch_log_processor_with_async_runtime", ] } -opentelemetry-otlp = { version = "0.32", features = ["trace", "metrics", "http-json"], default-features = false } +opentelemetry-otlp = { version = "0.32", features = ["trace", "metrics", "logs", "http-json"], default-features = false } opentelemetry-http = { version = "0.32", default-features = false } opentelemetry-semantic-conventions = "0.32" @@ -77,6 +82,14 @@ pin-project-lite = { version = "0.2", optional = true } # 0.5 supports reqwest 0.13 (this crate's pin). reqwest-middleware = { version = "0.5", optional = true } +# The STANDARD tracing→OTel-log bridge: a `tracing_subscriber::Layer` that turns +# every `tracing` event into an OTel `LogRecord` on our `SdkLoggerProvider`. The +# SDK logger enriches each record's trace_id/span_id from the ACTIVE +# `opentelemetry::Context` at emit time, so a log fired inside one of this crate's +# spans (tower/gen_ai/reqwest, all OTel-native) is correlated automatically. The +# host installs the layer via `OtelSdkHandle::tracing_appender_layer()`. +opentelemetry-appender-tracing = "0.32" + [features] default = [] # OTel server-span layer for Tower/Axum services. @@ -89,3 +102,9 @@ reqwest-middleware = ["dep:reqwest-middleware", "dep:reqwest"] [dev-dependencies] tokio = { version = "1", features = ["rt-multi-thread", "macros", "time", "test-util"] } wiremock = "0.6" +# Logs signal tests: emit `tracing` events through the bridge and assert on the +# captured records. `testing` unlocks `InMemoryLogExporter`; the feature unifies +# onto the normal `opentelemetry_sdk` dep only when building test targets. +tracing = "0.1" +tracing-subscriber = { version = "0.3", features = ["registry", "std"] } +opentelemetry_sdk = { version = "0.32", features = ["testing"] } diff --git a/rust/observability/src/bootstrap.rs b/rust/observability/src/bootstrap.rs index 40c752a..1a68687 100644 --- a/rust/observability/src/bootstrap.rs +++ b/rust/observability/src/bootstrap.rs @@ -9,8 +9,9 @@ //! ## Env vars (identical names to the TS SDK) //! //! - `SMOOAI_OBSERVABILITY_ENDPOINT` — base ingest URL (e.g. `https://api.smoo.ai`). -//! `/v1/traces` and `/v1/metrics` are appended. Per-signal overrides via the -//! standard `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` / `_METRICS_ENDPOINT`. +//! `/v1/traces`, `/v1/metrics`, and `/v1/logs` are appended. Per-signal +//! overrides via the standard `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` / +//! `_METRICS_ENDPOINT` / `_LOGS_ENDPOINT`. //! - Auth (pick one; pre-minted token wins): //! - `SMOOAI_OBSERVABILITY_TOKEN` — pre-minted Bearer JWT (not refreshed). //! - `SMOOAI_OBSERVABILITY_AUTH_URL` + `_CLIENT_ID` + `_CLIENT_SECRET` — @@ -36,6 +37,7 @@ pub struct BootstrapEnv { pub endpoint: Option, pub traces_endpoint: Option, pub metrics_endpoint: Option, + pub logs_endpoint: Option, pub token: Option, pub auth_url: Option, pub client_id: Option, @@ -54,6 +56,7 @@ impl BootstrapEnv { endpoint: env::var("SMOOAI_OBSERVABILITY_ENDPOINT").ok(), traces_endpoint: env::var("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").ok(), metrics_endpoint: env::var("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT").ok(), + logs_endpoint: env::var("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT").ok(), token: env::var("SMOOAI_OBSERVABILITY_TOKEN").ok(), auth_url: env::var("SMOOAI_OBSERVABILITY_AUTH_URL").ok(), client_id: env::var("SMOOAI_OBSERVABILITY_CLIENT_ID").ok(), @@ -127,6 +130,11 @@ async fn build(env: BootstrapEnv) -> BootstrapResult { .as_ref() .map(|e| format!("{}/v1/metrics", strip_trailing_slash(e))) }); + let logs_endpoint = env.logs_endpoint.clone().or_else(|| { + env.endpoint + .as_ref() + .map(|e| format!("{}/v1/logs", strip_trailing_slash(e))) + }); // Auth: static token wins; otherwise build a per-request TokenProvider and // warm it so the first export doesn't pay the round-trip. @@ -157,10 +165,12 @@ async fn build(env: BootstrapEnv) -> BootstrapResult { warn("no auth configured (set SMOOAI_OBSERVABILITY_TOKEN or _AUTH_URL/_CLIENT_ID/_CLIENT_SECRET); OTLP exports will be unauthenticated"); } - let otel = if traces_endpoint.is_some() || metrics_endpoint.is_some() { + let otel = if traces_endpoint.is_some() || metrics_endpoint.is_some() || logs_endpoint.is_some() + { let mut opts = SetupOtelOptions::new(service_name); opts.otlp_traces_endpoint = traces_endpoint; opts.otlp_metrics_endpoint = metrics_endpoint; + opts.otlp_logs_endpoint = logs_endpoint; opts.otlp_headers = static_headers; opts.environment = environment.clone(); opts.release = release.clone(); diff --git a/rust/observability/src/otel.rs b/rust/observability/src/otel.rs index 88f3635..b1e7790 100644 --- a/rust/observability/src/otel.rs +++ b/rust/observability/src/otel.rs @@ -27,9 +27,11 @@ use crate::auth::TokenProvider; use once_cell::sync::OnceCell; +use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; use opentelemetry_otlp::{ - MetricExporter, Protocol, SpanExporter, WithExportConfig, WithHttpConfig, + LogExporter, MetricExporter, Protocol, SpanExporter, WithExportConfig, WithHttpConfig, }; +use opentelemetry_sdk::logs::{SdkLogger, SdkLoggerProvider}; use opentelemetry_sdk::metrics::SdkMeterProvider; use opentelemetry_sdk::runtime; use opentelemetry_sdk::trace::SdkTracerProvider; @@ -37,6 +39,11 @@ use opentelemetry_sdk::Resource; use std::collections::HashMap; use std::time::Duration; +/// The tracing→OTel-log bridge layer type, concretely parameterized over this +/// crate's SDK logger provider. Hosts add it to their `tracing_subscriber` +/// registry to route `tracing` events into the OTLP logs pipeline. +pub type TracingAppenderLayer = OpenTelemetryTracingBridge; + mod auth_client; pub use auth_client::AuthInjectingHttpClient; @@ -50,6 +57,9 @@ pub struct SetupOtelOptions { /// Fully-qualified OTLP/HTTP endpoint for metrics. When `None`, metrics are /// not exported. pub otlp_metrics_endpoint: Option, + /// Fully-qualified OTLP/HTTP endpoint for logs (e.g. + /// `https://api.smoo.ai/v1/logs`). When `None`, logs are not exported. + pub otlp_logs_endpoint: Option, /// Static headers merged onto every export (e.g. a pre-minted /// `authorization` Bearer, or `user-agent`). pub otlp_headers: HashMap, @@ -70,6 +80,7 @@ impl SetupOtelOptions { service_name: service_name.into(), otlp_traces_endpoint: None, otlp_metrics_endpoint: None, + otlp_logs_endpoint: None, otlp_headers: HashMap::new(), environment: None, release: None, @@ -85,10 +96,11 @@ impl SetupOtelOptions { pub struct OtelSdkHandle { tracer_provider: Option, meter_provider: Option, + logger_provider: Option, } impl OtelSdkHandle { - /// Force-flush spans + metrics now. Best-effort; errors are swallowed. + /// Force-flush spans + metrics + logs now. Best-effort; errors are swallowed. pub fn flush(&self) { if let Some(tp) = &self.tracer_provider { let _ = tp.force_flush(); @@ -96,6 +108,9 @@ impl OtelSdkHandle { if let Some(mp) = &self.meter_provider { let _ = mp.force_flush(); } + if let Some(lp) = &self.logger_provider { + let _ = lp.force_flush(); + } } /// Graceful shutdown — drains and closes the pipelines. Best-effort. @@ -106,6 +121,29 @@ impl OtelSdkHandle { if let Some(mp) = &self.meter_provider { let _ = mp.shutdown(); } + if let Some(lp) = &self.logger_provider { + let _ = lp.shutdown(); + } + } + + /// The tracing→OTel-log bridge layer for this SDK's logs pipeline, or `None` + /// when logs are disabled (no logs endpoint configured). Add it to the host's + /// `tracing_subscriber` registry so `tracing` events become OTLP log records: + /// + /// ```ignore + /// use tracing_subscriber::prelude::*; + /// if let Some(layer) = handle.tracing_appender_layer() { + /// tracing_subscriber::registry().with(layer).init(); + /// } + /// ``` + /// + /// Records emitted inside an active OTel span (the tower/gen_ai/reqwest spans + /// this crate creates) carry that span's `trace_id`/`span_id` automatically — + /// the SDK logger reads `opentelemetry::Context::current()` at emit time. + pub fn tracing_appender_layer(&self) -> Option { + self.logger_provider + .as_ref() + .map(OpenTelemetryTracingBridge::new) } } @@ -190,7 +228,7 @@ fn build_and_install(options: SetupOtelOptions) -> OtelSdkHandle { .with_interval(options.metric_export_interval) .build(); let mp = SdkMeterProvider::builder() - .with_resource(resource) + .with_resource(resource.clone()) .with_reader(reader) .build(); opentelemetry::global::set_meter_provider(mp.clone()); @@ -199,9 +237,29 @@ fn build_and_install(options: SetupOtelOptions) -> OtelSdkHandle { None => None, }; + // Logs pipeline — same async-runtime processor rationale as traces/metrics + // above (SMOODEV-2045). Unlike traces/metrics there is no stable global logger + // provider to install; the provider is handed to the host via + // `OtelSdkHandle::tracing_appender_layer()`, which builds the tracing bridge + // from it. trace_id/span_id correlation is the SDK logger's job — it reads the + // active `opentelemetry::Context` at emit time. + let logger_provider = match build_log_exporter(&options) { + Some(exporter) => { + use opentelemetry_sdk::logs::log_processor_with_async_runtime::BatchLogProcessor; + let processor = BatchLogProcessor::builder(exporter, runtime::Tokio).build(); + let lp = SdkLoggerProvider::builder() + .with_resource(resource) + .with_log_processor(processor) + .build(); + Some(lp) + } + None => None, + }; + OtelSdkHandle { tracer_provider, meter_provider, + logger_provider, } } @@ -252,6 +310,25 @@ fn build_metric_exporter(options: &SetupOtelOptions) -> Option { } } +fn build_log_exporter(options: &SetupOtelOptions) -> Option { + let endpoint = options.otlp_logs_endpoint.clone()?; + let client = build_http_client(options); + let result = LogExporter::builder() + .with_http() + .with_protocol(Protocol::HttpJson) + .with_endpoint(endpoint) + .with_headers(options.otlp_headers.clone()) + .with_http_client(client) + .build(); + match result { + Ok(exporter) => Some(exporter), + Err(e) => { + warn(&format!("failed to build log exporter: {e}")); + None + } + } +} + pub(crate) fn warn(message: &str) { use std::io::Write; let _ = writeln!(std::io::stderr(), "[@smooai/observability/otel] {message}"); @@ -263,12 +340,30 @@ mod tests { #[test] fn no_endpoints_yields_disabled_handle() { - // A fresh handle built with no endpoints must have neither provider and - // must not panic. (We can't call the global-installing setup_otel_sdk in - // a unit test without polluting global state, so exercise build paths.) + // A fresh handle built with no endpoints must have no provider on any + // signal and must not panic. (We can't call the global-installing + // setup_otel_sdk in a unit test without polluting global state, so + // exercise build paths.) let opts = SetupOtelOptions::new("test-svc"); assert!(build_span_exporter(&opts).is_none()); assert!(build_metric_exporter(&opts).is_none()); + assert!(build_log_exporter(&opts).is_none()); + } + + #[test] + fn disabled_logs_pipeline_yields_no_appender_layer() { + // The logs pipeline must be a strict no-op when no logs endpoint is + // configured: no logger provider, so no tracing bridge for the host to + // install. Existing traces/metrics behavior is untouched. + let handle = OtelSdkHandle { + tracer_provider: None, + meter_provider: None, + logger_provider: None, + }; + assert!(handle.tracing_appender_layer().is_none()); + // Flush/shutdown on a fully-disabled handle are safe no-ops. + handle.flush(); + handle.shutdown(); } #[test] diff --git a/rust/observability/tests/logs_signal.rs b/rust/observability/tests/logs_signal.rs new file mode 100644 index 0000000..d8dd3ed --- /dev/null +++ b/rust/observability/tests/logs_signal.rs @@ -0,0 +1,105 @@ +//! Logs-signal integration tests (th-5dca7d / th-de3805). +//! +//! The whole point of the logs signal is trace↔log correlation: a log emitted +//! inside an active span must carry that span's real W3C `trace_id`/`span_id`, +//! so the observability product can join logs to traces. This test drives that +//! path deterministically with an in-memory exporter. +//! +//! It uses the exact bridge `OtelSdkHandle::tracing_appender_layer()` hands the +//! host — `OpenTelemetryTracingBridge::new(&logger_provider)` — and a +//! `simple` log processor (synchronous export on emit) so the assertion has no +//! flush-timing / runtime dependency. The correlation itself is the SDK logger's +//! doing: it reads `opentelemetry::Context::current()` at emit time, which is the +//! same OTel-native context our tower/gen_ai/reqwest spans make active. + +use opentelemetry::logs::{AnyValue, Severity}; +use opentelemetry::trace::{TraceContextExt, Tracer, TracerProvider}; +use opentelemetry::Context; +use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; +use opentelemetry_sdk::logs::{InMemoryLogExporter, SdkLoggerProvider}; +use opentelemetry_sdk::trace::SdkTracerProvider; +use tracing::info; +use tracing_subscriber::prelude::*; + +#[test] +fn log_within_active_span_carries_trace_and_span_id() { + // Logs pipeline: in-memory exporter behind a simple (synchronous) processor. + let exporter = InMemoryLogExporter::default(); + let logger_provider = SdkLoggerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + + // The bridge the host installs — identical to `tracing_appender_layer()`. + let bridge = OpenTelemetryTracingBridge::new(&logger_provider); + let subscriber = tracing_subscriber::registry().with(bridge); + + // A real sampled OTel span so its SpanContext carries valid W3C ids. + let tracer_provider = SdkTracerProvider::builder().build(); + let tracer = tracer_provider.tracer("logs-signal-test"); + let cx = Context::current().with_span(tracer.start("unit-of-work")); + let span_ctx = cx.span().span_context().clone(); + assert!( + span_ctx.is_valid(), + "test span must have a valid (sampled) span context" + ); + + tracing::subscriber::with_default(subscriber, || { + // Make the span the current OTel context, then log inside it. + let _guard = cx.clone().attach(); + info!("hello from within a span"); + }); + + logger_provider.force_flush().unwrap(); + + let logs = exporter.get_emitted_logs().expect("emitted logs readable"); + assert_eq!(logs.len(), 1, "exactly one record should have been emitted"); + let record = &logs[0].record; + + // Body ← message. + match record.body() { + Some(AnyValue::String(s)) => assert_eq!(s.as_str(), "hello from within a span"), + other => panic!("unexpected log body: {other:?}"), + } + + // Severity ← tracing level. + assert_eq!(record.severity_number(), Some(Severity::Info)); + + // trace_id / span_id ← ACTIVE span. This is the correlation the product needs. + let trace_context = record + .trace_context() + .expect("record emitted inside a span must carry trace context"); + assert_eq!( + trace_context.trace_id, + span_ctx.trace_id(), + "log record trace_id must match the active span's trace_id" + ); + assert_eq!( + trace_context.span_id, + span_ctx.span_id(), + "log record span_id must match the active span's span_id" + ); +} + +#[test] +fn log_without_active_span_has_no_trace_context() { + // Outside any span the record still flows through the pipeline, but carries + // no trace context — nothing to correlate to, and nothing fabricated. + let exporter = InMemoryLogExporter::default(); + let logger_provider = SdkLoggerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + let bridge = OpenTelemetryTracingBridge::new(&logger_provider); + let subscriber = tracing_subscriber::registry().with(bridge); + + tracing::subscriber::with_default(subscriber, || { + info!("no span here"); + }); + + logger_provider.force_flush().unwrap(); + let logs = exporter.get_emitted_logs().unwrap(); + assert_eq!(logs.len(), 1); + assert!( + logs[0].record.trace_context().is_none(), + "a log with no active span must not carry a (fabricated) trace context" + ); +}