Reload app-server OTel from thread config

This commit is contained in:
Eric Traut
2026-04-20 12:55:04 -07:00
parent cc96a03f10
commit d8621f0795
10 changed files with 331 additions and 14 deletions

View File

@@ -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<LogDbLayer>,
otel_reloader: Option<OtelReloader>,
}
#[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<OtelReloader>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -640,6 +643,7 @@ pub(crate) struct CodexMessageProcessorArgs {
pub(crate) cloud_requirements: Arc<RwLock<CloudRequirementsLoader>>,
pub(crate) feedback: CodexFeedback,
pub(crate) log_db: Option<LogDbLayer>,
pub(crate) otel_reloader: Option<OtelReloader>,
}
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 {

View File

@@ -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());

View File

@@ -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(())
}

View File

@@ -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<AuthManager>,
pub(crate) rpc_transport: AppServerRpcTransport,
pub(crate) remote_control_handle: Option<RemoteControlHandle>,
pub(crate) otel_reloader: Option<OtelReloader>,
}
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.

View File

@@ -252,6 +252,7 @@ fn build_test_processor(
auth_manager,
rpc_transport: AppServerRpcTransport::Stdio,
remote_control_handle: None,
otel_reloader: None,
}));
(processor, outgoing_rx)
}

View File

@@ -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<bool>,
default_analytics_enabled: bool,
}
struct OtelReloadState {
key: Option<OtelProviderKey>,
provider: Option<OtelProvider>,
}
#[derive(Clone)]
pub(crate) struct OtelReloader {
logger_layer: OtelLoggerLayer,
state: Arc<Mutex<OtelReloadState>>,
default_analytics_enabled: bool,
}
impl OtelReloader {
pub(crate) fn new(
initial_config: &Config,
provider: Option<OtelProvider>,
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<dyn Error>> {
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<dyn Error>> {
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::<String>()
);
Ok(())
}
async fn wait_for_otel_logs_payload(server: &MockServer) -> Result<Vec<u8>, Box<dyn Error>> {
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::<Vec<u8>, Box<dyn Error>>(request.body.clone());
}
tokio::time::sleep(Duration::from_millis(25)).await;
}
})
.await??;
Ok(body)
}
}

View File

@@ -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;

View File

@@ -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,
<SdkLoggerProvider as opentelemetry::logs::LoggerProvider>::Logger,
>;
#[derive(Clone, Default)]
pub struct OtelLoggerLayer {
bridge: Arc<RwLock<Option<LoggerBridge>>>,
}
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<S> Layer<S> 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<LoggerBridge> {
provider
.and_then(|provider| provider.logger.as_ref())
.map(OpenTelemetryTracingBridge::new)
}

View File

@@ -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<MetricsClient> = OnceLock::new();
static GLOBAL_METRICS: OnceLock<RwLock<Option<MetricsClient>>> = OnceLock::new();
pub(crate) fn install_global(metrics: MetricsClient) {
let _ = GLOBAL_METRICS.set(metrics);
fn global_metrics() -> &'static RwLock<Option<MetricsClient>> {
GLOBAL_METRICS.get_or_init(|| RwLock::new(None))
}
pub(crate) fn replace_global(metrics: Option<MetricsClient>) {
let mut global = global_metrics()
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner);
*global = metrics;
}
pub fn global() -> Option<MetricsClient> {
GLOBAL_METRICS.get().cloned()
global_metrics()
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}

View File

@@ -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.");