diff --git a/codex-rs/app-server/src/codex_message_processor.rs b/codex-rs/app-server/src/codex_message_processor.rs index 51b1c48c83..b1cb0074cb 100644 --- a/codex-rs/app-server/src/codex_message_processor.rs +++ b/codex-rs/app-server/src/codex_message_processor.rs @@ -11,6 +11,7 @@ use crate::fuzzy_file_search::FuzzyFileSearchSession; use crate::fuzzy_file_search::run_fuzzy_file_search; use crate::fuzzy_file_search::start_fuzzy_file_search_session; use crate::models::supported_models; +use crate::otel_reload::OtelReloader; use crate::outgoing_message::ConnectionId; use crate::outgoing_message::ConnectionRequestId; use crate::outgoing_message::OutgoingMessageSender; @@ -488,6 +489,7 @@ pub(crate) struct CodexMessageProcessor { background_tasks: TaskTracker, feedback: CodexFeedback, log_db: Option, + otel_reloader: Option, } #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] @@ -509,6 +511,7 @@ struct ListenerTaskContext { thread_watch_manager: ThreadWatchManager, fallback_model_provider: String, codex_home: PathBuf, + otel_reloader: Option, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -640,6 +643,7 @@ pub(crate) struct CodexMessageProcessorArgs { pub(crate) cloud_requirements: Arc>, pub(crate) feedback: CodexFeedback, pub(crate) log_db: Option, + pub(crate) otel_reloader: Option, } impl CodexMessageProcessor { @@ -718,6 +722,7 @@ impl CodexMessageProcessor { cloud_requirements, feedback, log_db, + otel_reloader, } = args; Self { auth_manager, @@ -740,6 +745,7 @@ impl CodexMessageProcessor { background_tasks: TaskTracker::new(), feedback, log_db, + otel_reloader, } } @@ -2420,6 +2426,7 @@ impl CodexMessageProcessor { thread_watch_manager: self.thread_watch_manager.clone(), fallback_model_provider: self.config.model_provider_id.clone(), codex_home: self.config.codex_home.to_path_buf(), + otel_reloader: self.otel_reloader.clone(), }; let request_trace = request_context.request_trace(); let runtime_feature_enablement = self.current_runtime_feature_enablement(); @@ -2621,6 +2628,27 @@ impl CodexMessageProcessor { }; } + if let Some(otel_reloader) = listener_task_context.otel_reloader.as_ref() { + let otel_reload_error = otel_reloader + .reload_from_config(&config) + .err() + .map(|err| err.to_string()); + if let Some(err) = otel_reload_error { + listener_task_context + .outgoing + .send_error( + request_id, + JSONRPCErrorError { + code: INTERNAL_ERROR_CODE, + message: format!("failed to reload otel config: {err}"), + data: None, + }, + ) + .await; + return; + } + } + let instruction_sources = Self::instruction_sources_from_config(&config).await; let dynamic_tools = dynamic_tools.unwrap_or_default(); let core_dynamic_tools = if dynamic_tools.is_empty() { @@ -7936,6 +7964,7 @@ impl CodexMessageProcessor { thread_watch_manager: self.thread_watch_manager.clone(), fallback_model_provider: self.config.model_provider_id.clone(), codex_home: self.config.codex_home.to_path_buf(), + otel_reloader: self.otel_reloader.clone(), }, conversation_id, connection_id, @@ -8050,6 +8079,7 @@ impl CodexMessageProcessor { thread_watch_manager: self.thread_watch_manager.clone(), fallback_model_provider: self.config.model_provider_id.clone(), codex_home: self.config.codex_home.to_path_buf(), + otel_reloader: self.otel_reloader.clone(), }, conversation_id, conversation, @@ -8099,6 +8129,7 @@ impl CodexMessageProcessor { thread_watch_manager, fallback_model_provider, codex_home, + otel_reloader: _, } = listener_task_context; let outgoing_for_task = Arc::clone(&outgoing); tokio::spawn(async move { diff --git a/codex-rs/app-server/src/in_process.rs b/codex-rs/app-server/src/in_process.rs index 4458bce89d..bcbf8a4ac7 100644 --- a/codex-rs/app-server/src/in_process.rs +++ b/codex-rs/app-server/src/in_process.rs @@ -404,6 +404,7 @@ fn start_uninitialized(args: InProcessStartArgs) -> InProcessClientHandle { auth_manager, rpc_transport: AppServerRpcTransport::InProcess, remote_control_handle: None, + otel_reloader: None, })); let mut thread_created_rx = processor.thread_created_receiver(); let session = Arc::new(ConnectionSessionState::default()); diff --git a/codex-rs/app-server/src/lib.rs b/codex-rs/app-server/src/lib.rs index b7b099960b..6bf9d1fb3b 100644 --- a/codex-rs/app-server/src/lib.rs +++ b/codex-rs/app-server/src/lib.rs @@ -80,6 +80,7 @@ mod fuzzy_file_search; pub mod in_process; mod message_processor; mod models; +mod otel_reload; mod outgoing_message; mod server_request_error; mod thread_state; @@ -516,15 +517,16 @@ pub async fn run_main_with_transport( let log_db_layer = log_db .clone() .map(|layer| layer.with_filter(Targets::new().with_default(Level::TRACE))); - let otel_logger_layer = otel.as_ref().and_then(|o| o.logger_layer()); - let otel_tracing_layer = otel.as_ref().and_then(|o| o.tracing_layer()); + let otel_tracing_layer = otel.as_ref().and_then(|provider| provider.tracing_layer()); + let (otel_layer, otel_reloader) = + otel_reload::OtelReloader::new(&config, otel, default_analytics_enabled); let _ = tracing_subscriber::registry() .with(stderr_fmt) + .with(otel_layer) + .with(otel_tracing_layer) .with(feedback_layer) .with(feedback_metadata_layer) .with(log_db_layer) - .with(otel_logger_layer) - .with(otel_tracing_layer) .try_init(); for warning in &config_warnings { match &warning.details { @@ -666,6 +668,7 @@ pub async fn run_main_with_transport( auth_manager, rpc_transport: analytics_rpc_transport(transport), remote_control_handle: Some(remote_control_handle), + otel_reloader: Some(otel_reloader.clone()), })); let mut thread_created_rx = processor.thread_created_receiver(); let mut running_turn_count_rx = processor.subscribe_running_assistant_turn_count(); @@ -887,9 +890,7 @@ pub async fn run_main_with_transport( let _ = handle.await; } - if let Some(otel) = otel { - otel.shutdown(); - } + otel_reloader.shutdown(); Ok(()) } diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index da72af81c5..a486cb84b6 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -14,6 +14,7 @@ use crate::error_code::INVALID_REQUEST_ERROR_CODE; use crate::external_agent_config_api::ExternalAgentConfigApi; use crate::fs_api::FsApi; use crate::fs_watch::FsWatchManager; +use crate::otel_reload::OtelReloader; use crate::outgoing_message::ConnectionId; use crate::outgoing_message::ConnectionRequestId; use crate::outgoing_message::OutgoingMessageSender; @@ -242,6 +243,7 @@ pub(crate) struct MessageProcessorArgs { pub(crate) auth_manager: Arc, pub(crate) rpc_transport: AppServerRpcTransport, pub(crate) remote_control_handle: Option, + pub(crate) otel_reloader: Option, } impl MessageProcessor { @@ -263,6 +265,7 @@ impl MessageProcessor { auth_manager, rpc_transport, remote_control_handle, + otel_reloader, } = args; auth_manager.set_external_auth(Arc::new(ExternalAuthRefreshBridge { outgoing: outgoing.clone(), @@ -303,6 +306,7 @@ impl MessageProcessor { cloud_requirements: cloud_requirements.clone(), feedback, log_db, + otel_reloader, }); // Keep plugin startup warmups aligned at app-server startup. // TODO(xl): Move into PluginManager once this no longer depends on config feature gating. diff --git a/codex-rs/app-server/src/message_processor/tracing_tests.rs b/codex-rs/app-server/src/message_processor/tracing_tests.rs index c1bb995bab..124d0efb7b 100644 --- a/codex-rs/app-server/src/message_processor/tracing_tests.rs +++ b/codex-rs/app-server/src/message_processor/tracing_tests.rs @@ -252,6 +252,7 @@ fn build_test_processor( auth_manager, rpc_transport: AppServerRpcTransport::Stdio, remote_control_handle: None, + otel_reloader: None, })); (processor, outgoing_rx) } diff --git a/codex-rs/app-server/src/otel_reload.rs b/codex-rs/app-server/src/otel_reload.rs new file mode 100644 index 0000000000..7404125318 --- /dev/null +++ b/codex-rs/app-server/src/otel_reload.rs @@ -0,0 +1,182 @@ +//! App-server OpenTelemetry reloading for configuration resolved after startup. +//! +//! The app server installs its tracing subscriber once, but `thread/start` can +//! later load project-scoped config from the requested cwd. This module keeps +//! the installed log layer stable while swapping the underlying OTel provider +//! when that effective thread config changes. + +use codex_config::types::OtelConfig; +use codex_core::config::Config; +use codex_otel::OtelLoggerLayer; +use codex_otel::OtelProvider; +use std::error::Error; +use std::sync::Arc; +use std::sync::Mutex; + +#[derive(Clone, Debug, PartialEq)] +struct OtelProviderKey { + otel: OtelConfig, + analytics_enabled: Option, + default_analytics_enabled: bool, +} + +struct OtelReloadState { + key: Option, + provider: Option, +} + +#[derive(Clone)] +pub(crate) struct OtelReloader { + logger_layer: OtelLoggerLayer, + state: Arc>, + default_analytics_enabled: bool, +} + +impl OtelReloader { + pub(crate) fn new( + initial_config: &Config, + provider: Option, + default_analytics_enabled: bool, + ) -> (OtelLoggerLayer, OtelReloader) { + let logger_layer = OtelLoggerLayer::from_provider(provider.as_ref()); + ( + logger_layer.clone(), + OtelReloader { + logger_layer, + state: Arc::new(Mutex::new(OtelReloadState { + key: Some(provider_key(initial_config, default_analytics_enabled)), + provider, + })), + default_analytics_enabled, + }, + ) + } + + pub(crate) fn reload_from_config(&self, config: &Config) -> Result<(), Box> { + let next_key = provider_key(config, self.default_analytics_enabled); + { + let state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if state.key.as_ref() == Some(&next_key) { + return Ok(()); + } + } + + let next_provider = codex_core::otel_init::build_provider( + config, + env!("CARGO_PKG_VERSION"), + Some("codex-app-server"), + self.default_analytics_enabled, + )?; + self.logger_layer.replace_provider(next_provider.as_ref()); + + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.key = Some(next_key); + state.provider = next_provider; + Ok(()) + } + + pub(crate) fn shutdown(&self) { + self.logger_layer.replace_provider(None); + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.provider.take(); + } +} + +fn provider_key(config: &Config, default_analytics_enabled: bool) -> OtelProviderKey { + OtelProviderKey { + otel: config.otel.clone(), + analytics_enabled: config.analytics_enabled, + default_analytics_enabled, + } +} + +#[cfg(test)] +mod tests { + use super::*; + use codex_config::types::OtelExporterKind; + use codex_config::types::OtelHttpProtocol; + use codex_core::config::ConfigBuilder; + use std::collections::HashMap; + use std::time::Duration; + use tempfile::TempDir; + use tokio::time::timeout; + use tracing_subscriber::prelude::*; + use wiremock::Mock; + use wiremock::MockServer; + use wiremock::ResponseTemplate; + use wiremock::matchers::method; + use wiremock::matchers::path; + + #[tokio::test(flavor = "multi_thread")] + async fn reload_from_config_updates_log_exporter() -> Result<(), Box> { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/v1/logs")) + .respond_with(ResponseTemplate::new(200)) + .mount(&server) + .await; + + let codex_home = TempDir::new()?; + let initial_config = ConfigBuilder::default() + .codex_home(codex_home.path().to_path_buf()) + .build() + .await?; + let (layer, reloader) = OtelReloader::new(&initial_config, /*provider*/ None, false); + let subscriber = tracing_subscriber::registry().with(layer); + let _guard = tracing::subscriber::set_default(subscriber); + + let mut project_config = initial_config.clone(); + project_config.otel.exporter = OtelExporterKind::OtlpHttp { + endpoint: format!("{}/v1/logs", server.uri()), + headers: HashMap::new(), + protocol: OtelHttpProtocol::Json, + tls: None, + }; + + reloader.reload_from_config(&project_config)?; + tracing::event!( + target: "codex_otel.log_only", + tracing::Level::INFO, + event.name = "codex.reload_test", + ); + reloader.shutdown(); + + let body = wait_for_otel_logs_payload(&server).await?; + let body = String::from_utf8(body)?; + assert!( + body.contains("codex.reload_test"), + "expected reloaded OTEL logs to include test event; body prefix: {}", + body.chars().take(2000).collect::() + ); + Ok(()) + } + + async fn wait_for_otel_logs_payload(server: &MockServer) -> Result, Box> { + let body = timeout(Duration::from_secs(10), async { + loop { + let Some(requests) = server.received_requests().await else { + tokio::time::sleep(Duration::from_millis(25)).await; + continue; + }; + if let Some(request) = requests + .iter() + .find(|request| request.method == "POST" && request.url.path() == "/v1/logs") + { + return Ok::, Box>(request.body.clone()); + } + tokio::time::sleep(Duration::from_millis(25)).await; + } + }) + .await??; + Ok(body) + } +} diff --git a/codex-rs/otel/src/lib.rs b/codex-rs/otel/src/lib.rs index c7d0b7c419..854768cd39 100644 --- a/codex-rs/otel/src/lib.rs +++ b/codex-rs/otel/src/lib.rs @@ -1,5 +1,6 @@ pub(crate) mod config; mod events; +mod logger_layer; pub(crate) mod metrics; pub(crate) mod provider; pub(crate) mod trace_context; @@ -18,6 +19,7 @@ pub use crate::config::OtelTlsConfig; pub use crate::events::session_telemetry::AuthEnvTelemetryMetadata; pub use crate::events::session_telemetry::SessionTelemetry; pub use crate::events::session_telemetry::SessionTelemetryMetadata; +pub use crate::logger_layer::OtelLoggerLayer; pub use crate::metrics::runtime_metrics::RuntimeMetricTotals; pub use crate::metrics::runtime_metrics::RuntimeMetricsSummary; pub use crate::metrics::timer::Timer; diff --git a/codex-rs/otel/src/logger_layer.rs b/codex-rs/otel/src/logger_layer.rs new file mode 100644 index 0000000000..0e4781210e --- /dev/null +++ b/codex-rs/otel/src/logger_layer.rs @@ -0,0 +1,86 @@ +//! Reloadable OpenTelemetry log layer. +//! +//! This layer stays registered with `tracing-subscriber` while allowing callers +//! to replace the OTel logger provider behind it. It is used by long-lived +//! processes that discover their effective telemetry config after subscriber +//! initialization. + +use crate::OtelProvider; +use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; +use opentelemetry_sdk::logs::SdkLoggerProvider; +use std::sync::Arc; +use std::sync::RwLock; +use tracing::subscriber::Interest; +use tracing_subscriber::Layer; +use tracing_subscriber::layer::Context; +use tracing_subscriber::registry::LookupSpan; + +type LoggerBridge = OpenTelemetryTracingBridge< + SdkLoggerProvider, + ::Logger, +>; + +#[derive(Clone, Default)] +pub struct OtelLoggerLayer { + bridge: Arc>>, +} + +impl OtelLoggerLayer { + pub fn from_provider(provider: Option<&OtelProvider>) -> Self { + Self { + bridge: Arc::new(RwLock::new(logger_bridge(provider))), + } + } + + pub fn replace_provider(&self, provider: Option<&OtelProvider>) { + let mut bridge = self + .bridge + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner); + *bridge = logger_bridge(provider); + } +} + +impl Layer for OtelLoggerLayer +where + S: tracing::Subscriber + for<'span> LookupSpan<'span>, +{ + fn register_callsite(&self, metadata: &'static tracing::Metadata<'static>) -> Interest { + if OtelProvider::log_export_filter(metadata) { + Interest::sometimes() + } else { + Interest::never() + } + } + + fn enabled(&self, metadata: &tracing::Metadata<'_>, _ctx: Context<'_, S>) -> bool { + if !OtelProvider::log_export_filter(metadata) { + return false; + } + + self.bridge + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .is_some() + } + + fn on_event(&self, event: &tracing::Event<'_>, ctx: Context<'_, S>) { + if !OtelProvider::log_export_filter(event.metadata()) { + return; + } + + let bridge = self + .bridge + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if let Some(bridge) = bridge.as_ref() { + bridge.on_event(event, ctx); + } + } +} + +fn logger_bridge(provider: Option<&OtelProvider>) -> Option { + provider + .and_then(|provider| provider.logger.as_ref()) + .map(OpenTelemetryTracingBridge::new) +} diff --git a/codex-rs/otel/src/metrics/mod.rs b/codex-rs/otel/src/metrics/mod.rs index bcbb85d35c..43bafb12a2 100644 --- a/codex-rs/otel/src/metrics/mod.rs +++ b/codex-rs/otel/src/metrics/mod.rs @@ -14,14 +14,25 @@ pub use crate::metrics::error::MetricsError; pub use crate::metrics::error::Result; pub use names::*; use std::sync::OnceLock; +use std::sync::RwLock; pub use tags::SessionMetricTagValues; -static GLOBAL_METRICS: OnceLock = OnceLock::new(); +static GLOBAL_METRICS: OnceLock>> = OnceLock::new(); -pub(crate) fn install_global(metrics: MetricsClient) { - let _ = GLOBAL_METRICS.set(metrics); +fn global_metrics() -> &'static RwLock> { + GLOBAL_METRICS.get_or_init(|| RwLock::new(None)) +} + +pub(crate) fn replace_global(metrics: Option) { + let mut global = global_metrics() + .write() + .unwrap_or_else(std::sync::PoisonError::into_inner); + *global = metrics; } pub fn global() -> Option { - GLOBAL_METRICS.get().cloned() + global_metrics() + .read() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .clone() } diff --git a/codex-rs/otel/src/provider.rs b/codex-rs/otel/src/provider.rs index b6df9f5ad2..7850fabdf9 100644 --- a/codex-rs/otel/src/provider.rs +++ b/codex-rs/otel/src/provider.rs @@ -84,9 +84,7 @@ impl OtelProvider { Some(MetricsClient::new(config)?) }; - if let Some(metrics) = metrics.as_ref() { - crate::metrics::install_global(metrics.clone()); - } + crate::metrics::replace_global(metrics.clone()); if !log_enabled && !trace_enabled && metrics.is_none() { debug!("No OTEL exporter enabled in settings.");