diff --git a/codex-rs/core/src/client.rs b/codex-rs/core/src/client.rs index ce354d9b0..642a6dec5 100644 --- a/codex-rs/core/src/client.rs +++ b/codex-rs/core/src/client.rs @@ -57,7 +57,7 @@ use codex_api::common::ResponsesWsRequest; use codex_api::create_text_param_for_request; use codex_api::error::ApiError; use codex_api::requests::responses::Compression; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::config_types::ReasoningSummary as ReasoningSummaryConfig; @@ -281,14 +281,14 @@ impl ModelClient { &self, prompt: &Prompt, model_info: &ModelInfo, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, ) -> Result> { if prompt.input.is_empty() { return Ok(Vec::new()); } let client_setup = self.current_client_setup().await?; let transport = ReqwestTransport::new(build_reqwest_client()); - let request_telemetry = Self::build_request_telemetry(otel_manager); + let request_telemetry = Self::build_request_telemetry(session_telemetry); let client = ApiCompactClient::new(transport, client_setup.api_provider, client_setup.api_auth) .with_telemetry(Some(request_telemetry)); @@ -321,7 +321,7 @@ impl ModelClient { raw_memories: Vec, model_info: &ModelInfo, effort: Option, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, ) -> Result> { if raw_memories.is_empty() { return Ok(Vec::new()); @@ -329,7 +329,7 @@ impl ModelClient { let client_setup = self.current_client_setup().await?; let transport = ReqwestTransport::new(build_reqwest_client()); - let request_telemetry = Self::build_request_telemetry(otel_manager); + let request_telemetry = Self::build_request_telemetry(session_telemetry); let client = ApiMemoriesClient::new(transport, client_setup.api_provider, client_setup.api_auth) .with_telemetry(Some(request_telemetry)); @@ -369,8 +369,8 @@ impl ModelClient { } /// Builds request telemetry for unary API calls (e.g., Compact endpoint). - fn build_request_telemetry(otel_manager: &OtelManager) -> Arc { - let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone())); + fn build_request_telemetry(session_telemetry: &SessionTelemetry) -> Arc { + let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone())); let request_telemetry: Arc = telemetry; request_telemetry } @@ -419,14 +419,14 @@ impl ModelClient { /// behavior remains consistent across both flows. async fn connect_websocket( &self, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, api_provider: codex_api::Provider, api_auth: CoreAuthProvider, turn_state: Option>>, turn_metadata_header: Option<&str>, ) -> std::result::Result { let headers = self.build_websocket_headers(turn_state.as_ref(), turn_metadata_header); - let websocket_telemetry = ModelClientSession::build_websocket_telemetry(otel_manager); + let websocket_telemetry = ModelClientSession::build_websocket_telemetry(session_telemetry); ApiWebSocketResponsesClient::new(api_provider, api_auth) .connect( headers, @@ -660,7 +660,7 @@ impl ModelClientSession { /// This performs only connection setup; it never sends prompt payloads. pub async fn preconnect_websocket( &mut self, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, model_info: &ModelInfo, ) -> std::result::Result<(), ApiError> { if !self.client.responses_websocket_enabled(model_info) { @@ -679,7 +679,7 @@ impl ModelClientSession { let connection = self .client .connect_websocket( - otel_manager, + session_telemetry, client_setup.api_provider, client_setup.api_auth, Some(Arc::clone(&self.turn_state)), @@ -692,7 +692,7 @@ impl ModelClientSession { /// Returns a websocket connection for this turn. async fn websocket_connection( &mut self, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, api_provider: codex_api::Provider, api_auth: CoreAuthProvider, turn_metadata_header: Option<&str>, @@ -713,7 +713,7 @@ impl ModelClientSession { let new_conn = self .client .connect_websocket( - otel_manager, + session_telemetry, api_provider, api_auth, Some(turn_state), @@ -751,7 +751,7 @@ impl ModelClientSession { &self, prompt: &Prompt, model_info: &ModelInfo, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, @@ -764,7 +764,7 @@ impl ModelClientSession { self.client.state.provider.stream_idle_timeout(), ) .map_err(map_api_error)?; - let (stream, _last_request_rx) = map_response_stream(stream, otel_manager.clone()); + let (stream, _last_request_rx) = map_response_stream(stream, session_telemetry.clone()); return Ok(stream); } @@ -775,7 +775,8 @@ impl ModelClientSession { loop { let client_setup = self.client.current_client_setup().await?; let transport = ReqwestTransport::new(build_reqwest_client()); - let (request_telemetry, sse_telemetry) = Self::build_streaming_telemetry(otel_manager); + let (request_telemetry, sse_telemetry) = + Self::build_streaming_telemetry(session_telemetry); let compression = self.responses_request_compression(client_setup.auth.as_ref()); let options = self.build_responses_options(turn_metadata_header, compression); @@ -797,7 +798,7 @@ impl ModelClientSession { match stream_result { Ok(stream) => { - let (stream, _) = map_response_stream(stream, otel_manager.clone()); + let (stream, _) = map_response_stream(stream, session_telemetry.clone()); return Ok(stream); } Err(ApiError::Transport( @@ -817,7 +818,7 @@ impl ModelClientSession { &mut self, prompt: &Prompt, model_info: &ModelInfo, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, @@ -852,7 +853,7 @@ impl ModelClientSession { match self .websocket_connection( - otel_manager, + session_telemetry, client_setup.api_provider, client_setup.api_auth, turn_metadata_header, @@ -890,7 +891,7 @@ impl ModelClientSession { .await .map_err(map_api_error)?; let (stream, last_request_rx) = - map_response_stream(stream_result, otel_manager.clone()); + map_response_stream(stream_result, session_telemetry.clone()); self.websocket_session.last_response_rx = Some(last_request_rx); return Ok(WebsocketStreamOutcome::Stream(stream)); } @@ -898,17 +899,19 @@ impl ModelClientSession { /// Builds request and SSE telemetry for streaming API calls. fn build_streaming_telemetry( - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, ) -> (Arc, Arc) { - let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone())); + let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone())); let request_telemetry: Arc = telemetry.clone(); let sse_telemetry: Arc = telemetry; (request_telemetry, sse_telemetry) } /// Builds telemetry for the Responses API WebSocket transport. - fn build_websocket_telemetry(otel_manager: &OtelManager) -> Arc { - let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone())); + fn build_websocket_telemetry( + session_telemetry: &SessionTelemetry, + ) -> Arc { + let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone())); let websocket_telemetry: Arc = telemetry; websocket_telemetry } @@ -918,7 +921,7 @@ impl ModelClientSession { &mut self, prompt: &Prompt, model_info: &ModelInfo, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, @@ -935,7 +938,7 @@ impl ModelClientSession { .stream_responses_websocket( prompt, model_info, - otel_manager, + session_telemetry, effort, summary, service_tier, @@ -956,7 +959,7 @@ impl ModelClientSession { Ok(()) } Ok(WebsocketStreamOutcome::FallbackToHttp) => { - self.try_switch_fallback_transport(otel_manager, model_info); + self.try_switch_fallback_transport(session_telemetry, model_info); Ok(()) } Err(err) => Err(err), @@ -974,7 +977,7 @@ impl ModelClientSession { &mut self, prompt: &Prompt, model_info: &ModelInfo, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, effort: Option, summary: ReasoningSummaryConfig, service_tier: Option, @@ -988,7 +991,7 @@ impl ModelClientSession { .stream_responses_websocket( prompt, model_info, - otel_manager, + session_telemetry, effort, summary, service_tier, @@ -999,7 +1002,7 @@ impl ModelClientSession { { WebsocketStreamOutcome::Stream(stream) => return Ok(stream), WebsocketStreamOutcome::FallbackToHttp => { - self.try_switch_fallback_transport(otel_manager, model_info); + self.try_switch_fallback_transport(session_telemetry, model_info); } } } @@ -1007,7 +1010,7 @@ impl ModelClientSession { self.stream_responses_api( prompt, model_info, - otel_manager, + session_telemetry, effort, summary, service_tier, @@ -1026,14 +1029,14 @@ impl ModelClientSession { /// Returns `true` if this call activated fallback, or `false` if fallback was already active. pub(crate) fn try_switch_fallback_transport( &mut self, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, model_info: &ModelInfo, ) -> bool { let websocket_enabled = self.client.responses_websocket_enabled(model_info); let activated = self.activate_http_fallback(websocket_enabled); if activated { warn!("falling back to HTTP"); - otel_manager.counter( + session_telemetry.counter( "codex.transport.fallback_to_http", 1, &[("from_wire_api", "responses_websocket")], @@ -1096,7 +1099,7 @@ fn build_responses_headers( fn map_response_stream( api_stream: S, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, ) -> (ResponseStream, oneshot::Receiver) where S: futures::Stream> @@ -1129,7 +1132,7 @@ where token_usage, }) => { if let Some(usage) = &token_usage { - otel_manager.sse_event_completed( + session_telemetry.sse_event_completed( usage.input_tokens, usage.output_tokens, Some(usage.cached_input_tokens), @@ -1162,7 +1165,7 @@ where Err(err) => { let mapped = map_api_error(err); if !logged_error { - otel_manager.see_event_completed_failed(&mapped); + session_telemetry.see_event_completed_failed(&mapped); logged_error = true; } if tx_event.send(Err(mapped)).await.is_err() { @@ -1198,12 +1201,12 @@ async fn handle_unauthorized( } struct ApiTelemetry { - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, } impl ApiTelemetry { - fn new(otel_manager: OtelManager) -> Self { - Self { otel_manager } + fn new(session_telemetry: SessionTelemetry) -> Self { + Self { session_telemetry } } } @@ -1216,7 +1219,7 @@ impl RequestTelemetry for ApiTelemetry { duration: Duration, ) { let error_message = error.map(std::string::ToString::to_string); - self.otel_manager.record_api_request( + self.session_telemetry.record_api_request( attempt, status.map(|s| s.as_u16()), error_message.as_deref(), @@ -1234,14 +1237,14 @@ impl SseTelemetry for ApiTelemetry { >, duration: Duration, ) { - self.otel_manager.log_sse_event(result, duration); + self.session_telemetry.log_sse_event(result, duration); } } impl WebsocketTelemetry for ApiTelemetry { fn on_ws_request(&self, duration: Duration, error: Option<&ApiError>) { let error_message = error.map(std::string::ToString::to_string); - self.otel_manager + self.session_telemetry .record_websocket_request(duration, error_message.as_deref()); } @@ -1250,14 +1253,15 @@ impl WebsocketTelemetry for ApiTelemetry { result: &std::result::Result>, ApiError>, duration: Duration, ) { - self.otel_manager.record_websocket_event(result, duration); + self.session_telemetry + .record_websocket_event(result, duration); } } #[cfg(test)] mod tests { use super::ModelClient; - use codex_otel::OtelManager; + use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::openai_models::ModelInfo; use codex_protocol::protocol::SessionSource; @@ -1313,8 +1317,8 @@ mod tests { .expect("deserialize test model info") } - fn test_otel_manager() -> OtelManager { - OtelManager::new( + fn test_session_telemetry() -> SessionTelemetry { + SessionTelemetry::new( ThreadId::new(), "gpt-test", "gpt-test", @@ -1344,10 +1348,10 @@ mod tests { async fn summarize_memories_returns_empty_for_empty_input() { let client = test_model_client(SessionSource::Cli); let model_info = test_model_info(); - let otel_manager = test_otel_manager(); + let session_telemetry = test_session_telemetry(); let output = client - .summarize_memories(Vec::new(), &model_info, None, &otel_manager) + .summarize_memories(Vec::new(), &model_info, None, &session_telemetry) .await .expect("empty summarize request should succeed"); assert_eq!(output.len(), 0); diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index bd32dfe02..519b28217 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -297,7 +297,7 @@ use crate::unified_exec::UnifiedExecProcessManager; use crate::util::backoff; use crate::windows_sandbox::WindowsSandboxLevelExt; use codex_async_utils::OrCancelExt; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_protocol::config_types::CollaborationMode; use codex_protocol::config_types::Personality; @@ -664,7 +664,7 @@ pub(crate) struct TurnContext { pub(crate) config: Arc, pub(crate) auth_manager: Option>, pub(crate) model_info: ModelInfo, - pub(crate) otel_manager: OtelManager, + pub(crate) session_telemetry: SessionTelemetry, pub(crate) provider: ModelProviderInfo, pub(crate) reasoning_effort: Option, pub(crate) reasoning_summary: ReasoningSummaryConfig, @@ -754,8 +754,8 @@ impl TurnContext { config: Arc::new(config), auth_manager: self.auth_manager.clone(), model_info: model_info.clone(), - otel_manager: self - .otel_manager + session_telemetry: self + .session_telemetry .clone() .with_model(model.as_str(), model_info.slug.as_str()), provider: self.provider.clone(), @@ -1089,7 +1089,7 @@ impl Session { #[allow(clippy::too_many_arguments)] fn make_turn_context( auth_manager: Option>, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, provider: ModelProviderInfo, session_configuration: &SessionConfiguration, per_turn_config: Config, @@ -1103,14 +1103,14 @@ impl Session { let reasoning_summary = session_configuration .model_reasoning_summary .unwrap_or(model_info.default_reasoning_summary); - let otel_manager = otel_manager.clone().with_model( + let session_telemetry = session_telemetry.clone().with_model( session_configuration.collaboration_mode.model(), model_info.slug.as_str(), ); let session_source = session_configuration.session_source.clone(); let auth_manager_for_context = auth_manager; let provider_for_context = provider; - let otel_manager_for_context = otel_manager; + let session_telemetry_for_context = session_telemetry; let per_turn_config = Arc::new(per_turn_config); let tools_config = ToolsConfig::new(&ToolsConfigParams { @@ -1140,7 +1140,7 @@ impl Session { config: per_turn_config.clone(), auth_manager: auth_manager_for_context, model_info: model_info.clone(), - otel_manager: otel_manager_for_context, + session_telemetry: session_telemetry_for_context, provider: provider_for_context, reasoning_effort, reasoning_summary, @@ -1355,7 +1355,7 @@ impl Session { let originator = crate::default_client::originator().value; let terminal_type = terminal::user_agent(); let session_model = session_configuration.collaboration_mode.model().to_string(); - let mut otel_manager = OtelManager::new( + let mut session_telemetry = SessionTelemetry::new( conversation_id, session_model.as_str(), session_model.as_str(), @@ -1368,7 +1368,7 @@ impl Session { session_configuration.session_source.clone(), ); if let Some(service_name) = session_configuration.metrics_service_name.as_deref() { - otel_manager = otel_manager.with_metrics_service_name(service_name); + session_telemetry = session_telemetry.with_metrics_service_name(service_name); } let network_proxy_audit_metadata = NetworkProxyAuditMetadata { conversation_id: Some(conversation_id.to_string()), @@ -1381,8 +1381,8 @@ impl Session { model: Some(session_model.clone()), slug: Some(session_model), }; - config.features.emit_metrics(&otel_manager); - otel_manager.counter( + config.features.emit_metrics(&session_telemetry); + session_telemetry.counter( "codex.thread.started", 1, &[( @@ -1395,7 +1395,7 @@ impl Session { )], ); - otel_manager.conversation_starts( + session_telemetry.conversation_starts( config.model_provider.name.as_str(), session_configuration.collaboration_mode.reasoning_effort(), config @@ -1438,7 +1438,7 @@ impl Session { conversation_id, session_configuration.cwd.clone(), &mut default_shell, - otel_manager.clone(), + session_telemetry.clone(), ) } } else { @@ -1533,7 +1533,7 @@ impl Session { show_raw_agent_reasoning: config.show_raw_agent_reasoning, exec_policy, auth_manager: Arc::clone(&auth_manager), - otel_manager, + session_telemetry, models_manager: Arc::clone(&models_manager), tool_approvals: Mutex::new(ApprovalStore::default()), execve_session_approvals: RwLock::new(HashMap::new()), @@ -2055,7 +2055,7 @@ impl Session { next_cwd.to_path_buf(), self.services.user_shell.as_ref().clone(), self.services.shell_snapshot_tx.clone(), - self.services.otel_manager.clone(), + self.services.session_telemetry.clone(), ); } @@ -2202,7 +2202,7 @@ impl Session { ); let mut turn_context: TurnContext = Self::make_turn_context( Some(Arc::clone(&self.services.auth_manager)), - &self.services.otel_manager, + &self.services.session_telemetry, session_configuration.provider.clone(), &session_configuration, per_turn_config, @@ -3024,7 +3024,7 @@ impl Session { pub(crate) async fn record_model_warning(&self, message: impl Into, ctx: &TurnContext) { self.services - .otel_manager + .session_telemetry .counter("codex.model_warning", 1, &[]); let item = ResponseItem::Message { id: None, @@ -4178,7 +4178,7 @@ mod handlers { }; sess.maybe_emit_unknown_model_warning_for_turn(current_context.as_ref()) .await; - current_context.otel_manager.user_prompt(&items); + current_context.session_telemetry.user_prompt(&items); // Attempt to inject input into current task. if let Err(SteerInputError::NoActiveTurn(items)) = sess.steer_input(items, None).await { @@ -4814,7 +4814,7 @@ mod handlers { .iter() .filter(|item| is_user_turn_boundary(item)) .count(); - sess.services.otel_manager.counter( + sess.services.session_telemetry.counter( "codex.conversation.turn.count", i64::try_from(turn_count).unwrap_or(0), &[], @@ -4933,13 +4933,13 @@ async fn spawn_review_thread( ); } - let otel_manager = parent_turn_context - .otel_manager + let session_telemetry = parent_turn_context + .session_telemetry .clone() .with_model(model.as_str(), review_model_info.slug.as_str()); let auth_manager_for_context = auth_manager.clone(); let provider_for_context = provider.clone(); - let otel_manager_for_context = otel_manager.clone(); + let session_telemetry_for_context = session_telemetry.clone(); let reasoning_effort = per_turn_config.model_reasoning_effort; let reasoning_summary = per_turn_config .model_reasoning_summary @@ -4965,7 +4965,7 @@ async fn spawn_review_thread( config: per_turn_config, auth_manager: auth_manager_for_context, model_info: model_info.clone(), - otel_manager: otel_manager_for_context, + session_telemetry: session_telemetry_for_context, provider: provider_for_context, reasoning_effort, reasoning_summary, @@ -5194,7 +5194,7 @@ pub(crate) async fn run_turn( ) .await; - let otel_manager = turn_context.otel_manager.clone(); + let session_telemetry = turn_context.session_telemetry.clone(); let thread_id = sess.conversation_id.to_string(); let tracking = build_track_events_context( turn_context.model_info.slug.clone(), @@ -5206,7 +5206,7 @@ pub(crate) async fn run_turn( warnings: skill_warnings, } = build_skill_injections( &mentioned_skills, - Some(&otel_manager), + Some(&session_telemetry), &sess.services.analytics_events_client, tracking.clone(), ) @@ -5829,8 +5829,10 @@ async fn run_sampling_request( // Use the configured provider-specific stream retry budget. let max_retries = turn_context.provider.stream_max_retries(); if retries >= max_retries - && client_session - .try_switch_fallback_transport(&turn_context.otel_manager, &turn_context.model_info) + && client_session.try_switch_fallback_transport( + &turn_context.session_telemetry, + &turn_context.model_info, + ) { sess.send_event( &turn_context, @@ -6533,7 +6535,7 @@ async fn try_run_sampling_request( .stream( prompt, &turn_context.model_info, - &turn_context.otel_manager, + &turn_context.session_telemetry, turn_context.reasoning_effort, turn_context.reasoning_summary, turn_context.config.service_tier, @@ -6589,7 +6591,7 @@ async fn try_run_sampling_request( }; sess.services - .otel_manager + .session_telemetry .record_responses(&handle_responses, &event); record_turn_ttft_metric(&turn_context, &event).await; diff --git a/codex-rs/core/src/codex_tests.rs b/codex-rs/core/src/codex_tests.rs index cb44d1cad..7fe6076f7 100644 --- a/codex-rs/core/src/codex_tests.rs +++ b/codex-rs/core/src/codex_tests.rs @@ -1811,13 +1811,13 @@ async fn build_test_config(codex_home: &Path) -> Config { .expect("load default test config") } -fn otel_manager( +fn session_telemetry( conversation_id: ThreadId, config: &Config, model_info: &ModelInfo, session_source: SessionSource, -) -> OtelManager { - OtelManager::new( +) -> SessionTelemetry { + SessionTelemetry::new( conversation_id, ModelsManager::get_model_offline_for_tests(config.model.as_deref()).as_str(), model_info.slug.as_str(), @@ -2026,7 +2026,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { session_configuration.collaboration_mode.model(), &per_turn_config, ); - let otel_manager = otel_manager( + let session_telemetry = session_telemetry( conversation_id, config.as_ref(), &model_info, @@ -2068,7 +2068,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { show_raw_agent_reasoning: config.show_raw_agent_reasoning, exec_policy, auth_manager: auth_manager.clone(), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), models_manager: Arc::clone(&models_manager), tool_approvals: Mutex::new(ApprovalStore::default()), execve_session_approvals: RwLock::new(HashMap::new()), @@ -2100,7 +2100,7 @@ pub(crate) async fn make_session_and_context() -> (Session, TurnContext) { let skills_outcome = Arc::new(services.skills_manager.skills_for_config(&per_turn_config)); let turn_context = Session::make_turn_context( Some(Arc::clone(&auth_manager)), - &otel_manager, + &session_telemetry, session_configuration.provider.clone(), &session_configuration, per_turn_config, @@ -2431,7 +2431,7 @@ pub(crate) async fn make_session_and_context_with_dynamic_tools_and_rx( session_configuration.collaboration_mode.model(), &per_turn_config, ); - let otel_manager = otel_manager( + let session_telemetry = session_telemetry( conversation_id, config.as_ref(), &model_info, @@ -2473,7 +2473,7 @@ pub(crate) async fn make_session_and_context_with_dynamic_tools_and_rx( show_raw_agent_reasoning: config.show_raw_agent_reasoning, exec_policy, auth_manager: Arc::clone(&auth_manager), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), models_manager: Arc::clone(&models_manager), tool_approvals: Mutex::new(ApprovalStore::default()), execve_session_approvals: RwLock::new(HashMap::new()), @@ -2505,7 +2505,7 @@ pub(crate) async fn make_session_and_context_with_dynamic_tools_and_rx( let skills_outcome = Arc::new(services.skills_manager.skills_for_config(&per_turn_config)); let turn_context = Arc::new(Session::make_turn_context( Some(Arc::clone(&auth_manager)), - &otel_manager, + &session_telemetry, session_configuration.provider.clone(), &session_configuration, per_turn_config, diff --git a/codex-rs/core/src/compact.rs b/codex-rs/core/src/compact.rs index 6c7dbb105..42e338443 100644 --- a/codex-rs/core/src/compact.rs +++ b/codex-rs/core/src/compact.rs @@ -400,7 +400,7 @@ async fn drain_to_completed( .stream( prompt, &turn_context.model_info, - &turn_context.otel_manager, + &turn_context.session_telemetry, turn_context.reasoning_effort, turn_context.reasoning_summary, turn_context.config.service_tier, diff --git a/codex-rs/core/src/compact_remote.rs b/codex-rs/core/src/compact_remote.rs index 40c520f65..473d91e54 100644 --- a/codex-rs/core/src/compact_remote.rs +++ b/codex-rs/core/src/compact_remote.rs @@ -107,7 +107,7 @@ async fn run_remote_compact_task_inner_impl( .compact_conversation_history( &prompt, &turn_context.model_info, - &turn_context.otel_manager, + &turn_context.session_telemetry, ) .or_else(|err| async { let total_usage_breakdown = sess.get_total_token_usage_breakdown().await; diff --git a/codex-rs/core/src/features.rs b/codex-rs/core/src/features.rs index 7defc571f..43c0b06b6 100644 --- a/codex-rs/core/src/features.rs +++ b/codex-rs/core/src/features.rs @@ -12,7 +12,7 @@ use crate::protocol::Event; use crate::protocol::EventMsg; use crate::protocol::WarningEvent; use codex_config::CONFIG_TOML_FILE; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use schemars::JsonSchema; use serde::Deserialize; use serde::Serialize; @@ -278,7 +278,7 @@ impl Features { self.legacy_usages.iter() } - pub fn emit_metrics(&self, otel: &OtelManager) { + pub fn emit_metrics(&self, otel: &SessionTelemetry) { for feature in FEATURES { if matches!(feature.stage, Stage::Removed) { continue; diff --git a/codex-rs/core/src/mcp_tool_call.rs b/codex-rs/core/src/mcp_tool_call.rs index 3900ea417..3dea0961c 100644 --- a/codex-rs/core/src/mcp_tool_call.rs +++ b/codex-rs/core/src/mcp_tool_call.rs @@ -107,7 +107,7 @@ pub(crate) async fn handle_mcp_tool_call( .await; let status = if result.is_ok() { "ok" } else { "error" }; turn_context - .otel_manager + .session_telemetry .counter("codex.mcp.call", 1, &[("status", status)]); return ResponseInputItem::McpToolCallOutput { call_id, result }; } @@ -188,7 +188,7 @@ pub(crate) async fn handle_mcp_tool_call( let status = if result.is_ok() { "ok" } else { "error" }; turn_context - .otel_manager + .session_telemetry .counter("codex.mcp.call", 1, &[("status", status)]); return ResponseInputItem::McpToolCallOutput { call_id, result }; @@ -229,7 +229,7 @@ pub(crate) async fn handle_mcp_tool_call( let status = if result.is_ok() { "ok" } else { "error" }; turn_context - .otel_manager + .session_telemetry .counter("codex.mcp.call", 1, &[("status", status)]); ResponseInputItem::McpToolCallOutput { call_id, result } diff --git a/codex-rs/core/src/memories/phase1.rs b/codex-rs/core/src/memories/phase1.rs index b860f84ee..ad4f29a0d 100644 --- a/codex-rs/core/src/memories/phase1.rs +++ b/codex-rs/core/src/memories/phase1.rs @@ -12,7 +12,7 @@ use crate::memories::prompts::build_stage_one_input_message; use crate::rollout::INTERACTIVE_SESSION_SOURCES; use crate::rollout::policy::should_persist_response_item_for_memories; use codex_api::ResponseEvent; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::config_types::ReasoningSummary as ReasoningSummaryConfig; use codex_protocol::config_types::ServiceTier; use codex_protocol::models::BaseInstructions; @@ -35,7 +35,7 @@ use tracing::warn; #[derive(Clone, Debug)] pub(in crate::memories) struct RequestContext { pub(in crate::memories) model_info: ModelInfo, - pub(in crate::memories) otel_manager: OtelManager, + pub(in crate::memories) session_telemetry: SessionTelemetry, pub(in crate::memories) reasoning_effort: Option, pub(in crate::memories) reasoning_summary: ReasoningSummaryConfig, pub(in crate::memories) service_tier: Option, @@ -85,7 +85,7 @@ struct StageOneOutput { pub(in crate::memories) async fn run(session: &Arc, config: &Config) { let _phase_one_e2e_timer = session .services - .otel_manager + .session_telemetry .start_timer(metrics::MEMORY_PHASE_ONE_E2E_MS, &[]) .ok(); @@ -94,7 +94,7 @@ pub(in crate::memories) async fn run(session: &Arc, config: &Config) { return; }; if claimed_candidates.is_empty() { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, 1, &[("status", "skipped_no_candidates")], @@ -168,7 +168,7 @@ impl RequestContext { Self { model_info, turn_metadata_header, - otel_manager: turn_context.otel_manager.clone(), + session_telemetry: turn_context.session_telemetry.clone(), reasoning_effort: Some(phase_one::REASONING_EFFORT), reasoning_summary: turn_context.reasoning_summary, service_tier: turn_context.config.service_tier, @@ -208,7 +208,7 @@ async fn claim_startup_jobs( Ok(claims) => Some(claims), Err(err) => { warn!("state db claim_stage1_jobs_for_startup failed during memories startup: {err}"); - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, 1, &[("status", "failed_claim")], @@ -347,7 +347,7 @@ mod job { .stream( &prompt, &stage_one_context.model_info, - &stage_one_context.otel_manager, + &stage_one_context.session_telemetry, stage_one_context.reasoning_effort, stage_one_context.reasoning_summary, stage_one_context.service_tier, @@ -516,60 +516,60 @@ fn aggregate_stats(outcomes: Vec) -> Stats { fn emit_metrics(session: &Session, counts: &Stats) { if counts.claimed > 0 { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, counts.claimed as i64, &[("status", "claimed")], ); } if counts.succeeded_with_output > 0 { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, counts.succeeded_with_output as i64, &[("status", "succeeded")], ); - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_OUTPUT, counts.succeeded_with_output as i64, &[], ); } if counts.succeeded_no_output > 0 { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, counts.succeeded_no_output as i64, &[("status", "succeeded_no_output")], ); } if counts.failed > 0 { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_ONE_JOBS, counts.failed as i64, &[("status", "failed")], ); } if let Some(token_usage) = counts.total_token_usage.as_ref() { - session.services.otel_manager.histogram( + session.services.session_telemetry.histogram( metrics::MEMORY_PHASE_ONE_TOKEN_USAGE, token_usage.total_tokens.max(0), &[("token_type", "total")], ); - session.services.otel_manager.histogram( + session.services.session_telemetry.histogram( metrics::MEMORY_PHASE_ONE_TOKEN_USAGE, token_usage.input_tokens.max(0), &[("token_type", "input")], ); - session.services.otel_manager.histogram( + session.services.session_telemetry.histogram( metrics::MEMORY_PHASE_ONE_TOKEN_USAGE, token_usage.cached_input(), &[("token_type", "cached_input")], ); - session.services.otel_manager.histogram( + session.services.session_telemetry.histogram( metrics::MEMORY_PHASE_ONE_TOKEN_USAGE, token_usage.output_tokens.max(0), &[("token_type", "output")], ); - session.services.otel_manager.histogram( + session.services.session_telemetry.histogram( metrics::MEMORY_PHASE_ONE_TOKEN_USAGE, token_usage.reasoning_output_tokens.max(0), &[("token_type", "reasoning_output")], diff --git a/codex-rs/core/src/memories/phase2.rs b/codex-rs/core/src/memories/phase2.rs index c8575035e..1a31bb335 100644 --- a/codex-rs/core/src/memories/phase2.rs +++ b/codex-rs/core/src/memories/phase2.rs @@ -43,7 +43,7 @@ struct Counters { pub(super) async fn run(session: &Arc, config: Arc) { let phase_two_e2e_timer = session .services - .otel_manager + .session_telemetry .start_timer(metrics::MEMORY_PHASE_TWO_E2E_MS, &[]) .ok(); @@ -59,7 +59,7 @@ pub(super) async fn run(session: &Arc, config: Arc) { let claim = match job::claim(session, db).await { Ok(claim) => claim, Err(e) => { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_TWO_JOBS, 1, &[("status", e)], @@ -183,7 +183,7 @@ mod job { session: &Arc, db: &StateRuntime, ) -> Result { - let otel_manager = &session.services.otel_manager; + let session_telemetry = &session.services.session_telemetry; let claim = db .try_claim_global_phase2_job(session.conversation_id, phase_two::JOB_LEASE_SECONDS) .await @@ -196,7 +196,11 @@ mod job { ownership_token, input_watermark, } => { - otel_manager.counter(metrics::MEMORY_PHASE_TWO_JOBS, 1, &[("status", "claimed")]); + session_telemetry.counter( + metrics::MEMORY_PHASE_TWO_JOBS, + 1, + &[("status", "claimed")], + ); (ownership_token, input_watermark) } codex_state::Phase2JobClaimOutcome::SkippedNotDirty => return Err("skipped_not_dirty"), @@ -212,7 +216,7 @@ mod job { claim: &Claim, reason: &'static str, ) { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_TWO_JOBS, 1, &[("status", reason)], @@ -244,7 +248,7 @@ mod job { selected_outputs: &[codex_state::Stage1Output], reason: &'static str, ) { - session.services.otel_manager.counter( + session.services.session_telemetry.counter( metrics::MEMORY_PHASE_TWO_JOBS, 1, &[("status", reason)], @@ -450,7 +454,7 @@ pub(super) fn get_watermark( } fn emit_metrics(session: &Arc, counters: Counters) { - let otel = session.services.otel_manager.clone(); + let otel = session.services.session_telemetry.clone(); if counters.input > 0 { otel.counter(metrics::MEMORY_PHASE_TWO_INPUT, counters.input, &[]); } @@ -463,7 +467,7 @@ fn emit_metrics(session: &Arc, counters: Counters) { } fn emit_token_usage_metrics(session: &Arc, token_usage: &TokenUsage) { - let otel = session.services.otel_manager.clone(); + let otel = session.services.session_telemetry.clone(); otel.histogram( metrics::MEMORY_PHASE_TWO_TOKEN_USAGE, token_usage.total_tokens.max(0), diff --git a/codex-rs/core/src/memories/usage.rs b/codex-rs/core/src/memories/usage.rs index ff82d6e16..8a86babb4 100644 --- a/codex-rs/core/src/memories/usage.rs +++ b/codex-rs/core/src/memories/usage.rs @@ -39,7 +39,7 @@ pub(crate) async fn emit_metric_for_tool_read(invocation: &ToolInvocation, succe let success = if success { "true" } else { "false" }; for kind in kinds { - invocation.turn.otel_manager.counter( + invocation.turn.session_telemetry.counter( MEMORIES_USAGE_METRIC, 1, &[ diff --git a/codex-rs/core/src/memory_trace.rs b/codex-rs/core/src/memory_trace.rs index 2e114b9aa..5cc499442 100644 --- a/codex-rs/core/src/memory_trace.rs +++ b/codex-rs/core/src/memory_trace.rs @@ -6,7 +6,7 @@ use crate::error::CodexErr; use crate::error::Result; use codex_api::RawMemory as ApiRawMemory; use codex_api::RawMemoryMetadata as ApiRawMemoryMetadata; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::openai_models::ModelInfo; use codex_protocol::openai_models::ReasoningEffort as ReasoningEffortConfig; use serde_json::Map; @@ -38,7 +38,7 @@ pub async fn build_memories_from_trace_files( trace_paths: &[PathBuf], model_info: &ModelInfo, effort: Option, - otel_manager: &OtelManager, + session_telemetry: &SessionTelemetry, ) -> Result> { if trace_paths.is_empty() { return Ok(Vec::new()); @@ -51,7 +51,7 @@ pub async fn build_memories_from_trace_files( let raw_memories = prepared.iter().map(|trace| trace.payload.clone()).collect(); let output = client - .summarize_memories(raw_memories, model_info, effort, otel_manager) + .summarize_memories(raw_memories, model_info, effort, session_telemetry) .await?; if output.len() != prepared.len() { return Err(CodexErr::InvalidRequest(format!( diff --git a/codex-rs/core/src/rollout/metadata.rs b/codex-rs/core/src/rollout/metadata.rs index 66ca9248e..3d46df0f8 100644 --- a/codex-rs/core/src/rollout/metadata.rs +++ b/codex-rs/core/src/rollout/metadata.rs @@ -7,7 +7,7 @@ use chrono::DateTime; use chrono::NaiveDateTime; use chrono::Timelike; use chrono::Utc; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::RolloutItem; @@ -96,7 +96,7 @@ pub(crate) fn builder_from_items( pub(crate) async fn extract_metadata_from_rollout( rollout_path: &Path, default_provider: &str, - otel: Option<&OtelManager>, + otel: Option<&SessionTelemetry>, ) -> anyhow::Result { let (items, _thread_id, parse_errors) = RolloutRecorder::load_rollout_items(rollout_path).await?; @@ -144,7 +144,7 @@ pub(crate) async fn extract_metadata_from_rollout( pub(crate) async fn backfill_sessions( runtime: &codex_state::StateRuntime, config: &Config, - otel: Option<&OtelManager>, + otel: Option<&SessionTelemetry>, ) { let timer = otel.and_then(|otel| otel.start_timer(DB_METRIC_BACKFILL_DURATION_MS, &[]).ok()); let backfill_state = match runtime.get_backfill_state().await { diff --git a/codex-rs/core/src/shell_snapshot.rs b/codex-rs/core/src/shell_snapshot.rs index 291129f97..5d285cbfa 100644 --- a/codex-rs/core/src/shell_snapshot.rs +++ b/codex-rs/core/src/shell_snapshot.rs @@ -14,7 +14,7 @@ use anyhow::Context; use anyhow::Result; use anyhow::anyhow; use anyhow::bail; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use tokio::fs; use tokio::process::Command; @@ -40,7 +40,7 @@ impl ShellSnapshot { session_id: ThreadId, session_cwd: PathBuf, shell: &mut Shell, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, ) -> watch::Sender>> { let (shell_snapshot_tx, shell_snapshot_rx) = watch::channel(None); shell.shell_snapshot = shell_snapshot_rx; @@ -51,7 +51,7 @@ impl ShellSnapshot { session_cwd, shell.clone(), shell_snapshot_tx.clone(), - otel_manager, + session_telemetry, ); shell_snapshot_tx @@ -63,7 +63,7 @@ impl ShellSnapshot { session_cwd: PathBuf, shell: Shell, shell_snapshot_tx: watch::Sender>>, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, ) { Self::spawn_snapshot_task( codex_home, @@ -71,7 +71,7 @@ impl ShellSnapshot { session_cwd, shell, shell_snapshot_tx, - otel_manager, + session_telemetry, ); } @@ -81,12 +81,12 @@ impl ShellSnapshot { session_cwd: PathBuf, snapshot_shell: Shell, shell_snapshot_tx: watch::Sender>>, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, ) { let snapshot_span = info_span!("shell_snapshot", thread_id = %session_id); tokio::spawn( async move { - let timer = otel_manager.start_timer("codex.shell_snapshot.duration_ms", &[]); + let timer = session_telemetry.start_timer("codex.shell_snapshot.duration_ms", &[]); let snapshot = ShellSnapshot::try_new( &codex_home, session_id, @@ -102,7 +102,7 @@ impl ShellSnapshot { if let Some(failure_reason) = snapshot.as_ref().err() { counter_tags.push(("failure_reason", *failure_reason)); } - otel_manager.counter("codex.shell_snapshot", 1, &counter_tags); + session_telemetry.counter("codex.shell_snapshot", 1, &counter_tags); let _ = shell_snapshot_tx.send(snapshot.ok()); } .instrument(snapshot_span), diff --git a/codex-rs/core/src/skills/injection.rs b/codex-rs/core/src/skills/injection.rs index de50d6cdd..ac55f3b8f 100644 --- a/codex-rs/core/src/skills/injection.rs +++ b/codex-rs/core/src/skills/injection.rs @@ -9,7 +9,7 @@ use crate::analytics_client::TrackEventsContext; use crate::instructions::SkillInstructions; use crate::mentions::build_skill_name_counts; use crate::skills::SkillMetadata; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::models::ResponseItem; use codex_protocol::user_input::UserInput; use tokio::fs; @@ -22,7 +22,7 @@ pub(crate) struct SkillInjections { pub(crate) async fn build_skill_injections( mentioned_skills: &[SkillMetadata], - otel: Option<&OtelManager>, + otel: Option<&SessionTelemetry>, analytics_client: &AnalyticsEventsClient, tracking: TrackEventsContext, ) -> SkillInjections { @@ -69,7 +69,11 @@ pub(crate) async fn build_skill_injections( result } -fn emit_skill_injected_metric(otel: Option<&OtelManager>, skill: &SkillMetadata, status: &str) { +fn emit_skill_injected_metric( + otel: Option<&SessionTelemetry>, + skill: &SkillMetadata, + status: &str, +) { let Some(otel) = otel else { return; }; diff --git a/codex-rs/core/src/skills/invocation_utils.rs b/codex-rs/core/src/skills/invocation_utils.rs index 36bb0c3ca..c4310bace 100644 --- a/codex-rs/core/src/skills/invocation_utils.rs +++ b/codex-rs/core/src/skills/invocation_utils.rs @@ -94,7 +94,7 @@ pub(crate) async fn maybe_emit_implicit_skill_invocation( return; } - turn_context.otel_manager.counter( + turn_context.session_telemetry.counter( "codex.skill.injected", 1, &[ diff --git a/codex-rs/core/src/state/service.rs b/codex-rs/core/src/state/service.rs index e2792638e..012e17bbf 100644 --- a/codex-rs/core/src/state/service.rs +++ b/codex-rs/core/src/state/service.rs @@ -20,7 +20,7 @@ use crate::tools::runtimes::ExecveSessionApproval; use crate::tools::sandboxing::ApprovalStore; use crate::unified_exec::UnifiedExecProcessManager; use codex_hooks::Hooks; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_utils_absolute_path::AbsolutePathBuf; use std::path::PathBuf; use tokio::sync::Mutex; @@ -45,7 +45,7 @@ pub(crate) struct SessionServices { pub(crate) exec_policy: ExecPolicyManager, pub(crate) auth_manager: Arc, pub(crate) models_manager: Arc, - pub(crate) otel_manager: OtelManager, + pub(crate) session_telemetry: SessionTelemetry, pub(crate) tool_approvals: Mutex, #[cfg_attr(not(unix), allow(dead_code))] pub(crate) execve_session_approvals: RwLock>, diff --git a/codex-rs/core/src/state_db.rs b/codex-rs/core/src/state_db.rs index 2a9fa5e45..922331e3c 100644 --- a/codex-rs/core/src/state_db.rs +++ b/codex-rs/core/src/state_db.rs @@ -7,7 +7,7 @@ use chrono::DateTime; use chrono::NaiveDateTime; use chrono::Timelike; use chrono::Utc; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::dynamic_tools::DynamicToolSpec; use codex_protocol::protocol::RolloutItem; @@ -26,7 +26,10 @@ pub type StateDbHandle = Arc; /// Initialize the state runtime for thread state persistence and backfill checks. To only be used /// inside `core`. The initialization should not be done anywhere else. -pub(crate) async fn init(config: &Config, otel: Option<&OtelManager>) -> Option { +pub(crate) async fn init( + config: &Config, + otel: Option<&SessionTelemetry>, +) -> Option { let runtime = match codex_state::StateRuntime::init( config.sqlite_home.clone(), config.model_provider_id.clone(), @@ -69,7 +72,10 @@ pub(crate) async fn init(config: &Config, otel: Option<&OtelManager>) -> Option< } /// Get the DB if the feature is enabled and the DB exists. -pub async fn get_state_db(config: &Config, otel: Option<&OtelManager>) -> Option { +pub async fn get_state_db( + config: &Config, + otel: Option<&SessionTelemetry>, +) -> Option { let state_path = codex_state::state_db_path(config.sqlite_home.as_path()); if !tokio::fs::try_exists(&state_path).await.unwrap_or(false) { return None; diff --git a/codex-rs/core/src/stream_events_utils.rs b/codex-rs/core/src/stream_events_utils.rs index 7f2ddb251..106f88105 100644 --- a/codex-rs/core/src/stream_events_utils.rs +++ b/codex-rs/core/src/stream_events_utils.rs @@ -220,7 +220,7 @@ pub(crate) async fn handle_output_item_done( Err(FunctionCallError::MissingLocalShellCallId) => { let msg = "LocalShellCall without call_id or id"; ctx.turn_context - .otel_manager + .session_telemetry .log_tool_failed("local_shell", msg); tracing::error!(msg); diff --git a/codex-rs/core/src/tasks/compact.rs b/codex-rs/core/src/tasks/compact.rs index a8c0a5e06..be2fcca85 100644 --- a/codex-rs/core/src/tasks/compact.rs +++ b/codex-rs/core/src/tasks/compact.rs @@ -30,14 +30,14 @@ impl SessionTask for CompactTask { ) -> Option { let session = session.clone_session(); let _ = if crate::compact::should_use_remote_compact_task(&ctx.provider) { - let _ = session.services.otel_manager.counter( + let _ = session.services.session_telemetry.counter( "codex.task.compact", 1, &[("type", "remote")], ); crate::compact_remote::run_remote_compact_task(session.clone(), ctx).await } else { - let _ = session.services.otel_manager.counter( + let _ = session.services.session_telemetry.counter( "codex.task.compact", 1, &[("type", "local")], diff --git a/codex-rs/core/src/tasks/mod.rs b/codex-rs/core/src/tasks/mod.rs index f6962edfd..a07e52f73 100644 --- a/codex-rs/core/src/tasks/mod.rs +++ b/codex-rs/core/src/tasks/mod.rs @@ -144,7 +144,7 @@ impl Session { let done = Arc::new(Notify::new()); let timer = turn_context - .otel_manager + .session_telemetry .start_timer("codex.turn.e2e_duration_ms", &[]) .ok(); @@ -272,7 +272,7 @@ impl Session { "false" }, ); - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.tool.call", i64::try_from(turn_tool_calls).unwrap_or(i64::MAX), &[tmp_mem], @@ -295,27 +295,27 @@ impl Session { - token_usage_at_turn_start.total_tokens) .max(0), }; - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.token_usage", turn_token_usage.total_tokens, &[("token_type", "total"), tmp_mem], ); - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.token_usage", turn_token_usage.input_tokens, &[("token_type", "input"), tmp_mem], ); - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.token_usage", turn_token_usage.cached_input(), &[("token_type", "cached_input"), tmp_mem], ); - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.token_usage", turn_token_usage.output_tokens, &[("token_type", "output"), tmp_mem], ); - self.services.otel_manager.histogram( + self.services.session_telemetry.histogram( "codex.turn.token_usage", turn_token_usage.reasoning_output_tokens, &[("token_type", "reasoning_output"), tmp_mem], diff --git a/codex-rs/core/src/tasks/regular.rs b/codex-rs/core/src/tasks/regular.rs index e3a9ef4c3..78bbd1db4 100644 --- a/codex-rs/core/src/tasks/regular.rs +++ b/codex-rs/core/src/tasks/regular.rs @@ -41,7 +41,7 @@ impl RegularTask { .prewarm_websocket( &prompt, &turn_context.model_info, - &turn_context.otel_manager, + &turn_context.session_telemetry, turn_context.reasoning_effort, turn_context.reasoning_summary, turn_context.config.service_tier, diff --git a/codex-rs/core/src/tasks/review.rs b/codex-rs/core/src/tasks/review.rs index 4787ee782..51ac3a1d3 100644 --- a/codex-rs/core/src/tasks/review.rs +++ b/codex-rs/core/src/tasks/review.rs @@ -57,7 +57,7 @@ impl SessionTask for ReviewTask { let _ = session .session .services - .otel_manager + .session_telemetry .counter("codex.task.review", 1, &[]); // Start sub-codex conversation and get the receiver for events. diff --git a/codex-rs/core/src/tasks/undo.rs b/codex-rs/core/src/tasks/undo.rs index 95f51c1ec..05cd928b5 100644 --- a/codex-rs/core/src/tasks/undo.rs +++ b/codex-rs/core/src/tasks/undo.rs @@ -45,7 +45,7 @@ impl SessionTask for UndoTask { let _ = session .session .services - .otel_manager + .session_telemetry .counter("codex.task.undo", 1, &[]); let sess = session.clone_session(); sess.send_event( diff --git a/codex-rs/core/src/tasks/user_shell.rs b/codex-rs/core/src/tasks/user_shell.rs index 2f77d9fce..a6a1301cd 100644 --- a/codex-rs/core/src/tasks/user_shell.rs +++ b/codex-rs/core/src/tasks/user_shell.rs @@ -98,7 +98,7 @@ pub(crate) async fn execute_user_shell_command( ) { session .services - .otel_manager + .session_telemetry .counter("codex.task.user_shell", 1, &[]); if mode == UserShellCommandMode::StandaloneTurn { diff --git a/codex-rs/core/src/tools/handlers/multi_agents.rs b/codex-rs/core/src/tools/handlers/multi_agents.rs index 20b001315..1d3805b32 100644 --- a/codex-rs/core/src/tools/handlers/multi_agents.rs +++ b/codex-rs/core/src/tools/handlers/multi_agents.rs @@ -214,7 +214,7 @@ mod spawn { .await; let new_thread_id = result?; let role_tag = role_name.unwrap_or(DEFAULT_ROLE_NAME); - turn.otel_manager + turn.session_telemetry .counter("codex.multi_agent.spawn", 1, &[("role", role_tag)]); let content = serde_json::to_string(&SpawnAgentResult { @@ -425,7 +425,7 @@ mod resume_agent { if let Some(err) = error { return Err(err); } - turn.otel_manager + turn.session_telemetry .counter("codex.multi_agent.resume", 1, &[]); let content = serde_json::to_string(&ResumeAgentResult { status }).map_err(|err| { diff --git a/codex-rs/core/src/tools/orchestrator.rs b/codex-rs/core/src/tools/orchestrator.rs index f66c79bbb..e78f0bcfc 100644 --- a/codex-rs/core/src/tools/orchestrator.rs +++ b/codex-rs/core/src/tools/orchestrator.rs @@ -108,7 +108,7 @@ impl ToolOrchestrator { where T: ToolRuntime, { - let otel = turn_ctx.otel_manager.clone(); + let otel = turn_ctx.session_telemetry.clone(); let otel_tn = &tool_ctx.tool_name; let otel_ci = &tool_ctx.call_id; let otel_user = ToolDecisionSource::User; diff --git a/codex-rs/core/src/tools/registry.rs b/codex-rs/core/src/tools/registry.rs index 77d3c4150..74a44f1db 100644 --- a/codex-rs/core/src/tools/registry.rs +++ b/codex-rs/core/src/tools/registry.rs @@ -82,7 +82,7 @@ impl ToolRegistry { ) -> Result { let tool_name = invocation.tool_name.clone(); let call_id_owned = invocation.call_id.clone(); - let otel = invocation.turn.otel_manager.clone(); + let otel = invocation.turn.session_telemetry.clone(); let payload_for_response = invocation.payload.clone(); let log_payload = payload_for_response.log_payload(); let metric_tags = [ diff --git a/codex-rs/core/src/tools/sandboxing.rs b/codex-rs/core/src/tools/sandboxing.rs index 28d87b5bf..df8f7c152 100644 --- a/codex-rs/core/src/tools/sandboxing.rs +++ b/codex-rs/core/src/tools/sandboxing.rs @@ -88,7 +88,7 @@ where let decision = fetch().await; - services.otel_manager.counter( + services.session_telemetry.counter( "codex.approval.requested", 1, &[ diff --git a/codex-rs/core/src/turn_timing.rs b/codex-rs/core/src/turn_timing.rs index b825fde94..765373c35 100644 --- a/codex-rs/core/src/turn_timing.rs +++ b/codex-rs/core/src/turn_timing.rs @@ -21,7 +21,7 @@ pub(crate) async fn record_turn_ttft_metric(turn_context: &TurnContext, event: & return; }; turn_context - .otel_manager + .session_telemetry .record_duration(TURN_TTFT_DURATION_METRIC, duration, &[]); } @@ -34,7 +34,7 @@ pub(crate) async fn record_turn_ttfm_metric(turn_context: &TurnContext, item: &T return; }; turn_context - .otel_manager + .session_telemetry .record_duration(TURN_TTFM_DURATION_METRIC, duration, &[]); } diff --git a/codex-rs/core/tests/responses_headers.rs b/codex-rs/core/tests/responses_headers.rs index be41a25a0..d5376fcd7 100644 --- a/codex-rs/core/tests/responses_headers.rs +++ b/codex-rs/core/tests/responses_headers.rs @@ -7,7 +7,7 @@ use codex_core::ModelProviderInfo; use codex_core::Prompt; use codex_core::ResponseEvent; use codex_core::WireApi; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_protocol::ThreadId; use codex_protocol::config_types::ReasoningSummary; @@ -72,7 +72,7 @@ async fn responses_stream_includes_subagent_header_on_review() { let session_source = SessionSource::SubAgent(SubAgentSource::Review); let model_info = codex_core::test_support::construct_model_info_offline(model.as_str(), &config); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( conversation_id, model.as_str(), model_info.slug.as_str(), @@ -113,7 +113,7 @@ async fn responses_stream_includes_subagent_header_on_review() { .stream( &prompt, &model_info, - &otel_manager, + &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), None, @@ -185,7 +185,7 @@ async fn responses_stream_includes_subagent_header_on_other() { let model_info = codex_core::test_support::construct_model_info_offline(model.as_str(), &config); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( conversation_id, model.as_str(), model_info.slug.as_str(), @@ -226,7 +226,7 @@ async fn responses_stream_includes_subagent_header_on_other() { .stream( &prompt, &model_info, - &otel_manager, + &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), None, @@ -297,7 +297,7 @@ async fn responses_respects_model_info_overrides_from_config() { SessionSource::SubAgent(SubAgentSource::Other("override-check".to_string())); let model_info = codex_core::test_support::construct_model_info_offline(model.as_str(), &config); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( conversation_id, model.as_str(), model_info.slug.as_str(), @@ -338,7 +338,7 @@ async fn responses_respects_model_info_overrides_from_config() { .stream( &prompt, &model_info, - &otel_manager, + &session_telemetry, effort, summary.unwrap_or(model_info.default_reasoning_summary), None, diff --git a/codex-rs/core/tests/suite/client.rs b/codex-rs/core/tests/suite/client.rs index 3ac9a72c3..30eac507c 100644 --- a/codex-rs/core/tests/suite/client.rs +++ b/codex-rs/core/tests/suite/client.rs @@ -12,7 +12,7 @@ use codex_core::default_client::originator; use codex_core::error::CodexErr; use codex_core::features::Feature; use codex_core::models_manager::collaboration_mode_presets::CollaborationModesConfig; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_protocol::ThreadId; use codex_protocol::config_types::CollaborationMode; @@ -1752,7 +1752,7 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() { let conversation_id = ThreadId::new(); let auth_manager = codex_core::test_support::auth_manager_from_auth(CodexAuth::from_api_key("Test API Key")); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( conversation_id, model.as_str(), model_info.slug.as_str(), @@ -1844,7 +1844,7 @@ async fn azure_responses_request_includes_store_and_reasoning_ids() { .stream( &prompt, &model_info, - &otel_manager, + &session_telemetry, effort, summary.unwrap_or(ReasoningSummary::Auto), None, diff --git a/codex-rs/core/tests/suite/client_websockets.rs b/codex-rs/core/tests/suite/client_websockets.rs index 4896a0fa1..cda634448 100755 --- a/codex-rs/core/tests/suite/client_websockets.rs +++ b/codex-rs/core/tests/suite/client_websockets.rs @@ -9,7 +9,7 @@ use codex_core::WireApi; use codex_core::X_RESPONSESAPI_INCLUDE_TIMING_METRICS_HEADER; use codex_core::features::Feature; use codex_core::ws_version_from_features; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_otel::metrics::MetricsClient; use codex_otel::metrics::MetricsConfig; @@ -56,7 +56,7 @@ struct WebsocketTestHarness { model_info: ModelInfo, effort: Option, summary: ReasoningSummary, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -105,7 +105,7 @@ async fn responses_websocket_preconnect_reuses_connection() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); client_session - .preconnect_websocket(&harness.otel_manager, &harness.model_info) + .preconnect_websocket(&harness.session_telemetry, &harness.model_info) .await .expect("websocket preconnect failed"); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -134,7 +134,7 @@ async fn responses_websocket_request_prewarm_reuses_connection() { .prewarm_websocket( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -207,7 +207,7 @@ async fn responses_websocket_preconnect_is_reused_even_with_header_changes() { let harness = websocket_harness(&server).await; let mut client_session = harness.client.new_session(); client_session - .preconnect_websocket(&harness.otel_manager, &harness.model_info) + .preconnect_websocket(&harness.session_telemetry, &harness.model_info) .await .expect("websocket preconnect failed"); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -215,7 +215,7 @@ async fn responses_websocket_preconnect_is_reused_even_with_header_changes() { .stream( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -253,7 +253,7 @@ async fn responses_websocket_request_prewarm_is_reused_even_with_header_changes( .prewarm_websocket( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -265,7 +265,7 @@ async fn responses_websocket_request_prewarm_is_reused_even_with_header_changes( .stream( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -318,7 +318,7 @@ async fn responses_websocket_prewarm_uses_v2_when_model_prefers_websockets_and_f .prewarm_websocket( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -371,7 +371,7 @@ async fn responses_websocket_preconnect_runs_when_only_v2_feature_enabled() { let harness = websocket_harness_with_options(&server, false, false, true, false).await; let mut client_session = harness.client.new_session(); client_session - .preconnect_websocket(&harness.otel_manager, &harness.model_info) + .preconnect_websocket(&harness.session_telemetry, &harness.model_info) .await .expect("websocket preconnect failed"); @@ -551,7 +551,7 @@ async fn responses_websocket_emits_websocket_telemetry_events() { .await; let harness = websocket_harness(&server).await; - harness.otel_manager.reset_runtime_metrics(); + harness.session_telemetry.reset_runtime_metrics(); let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -560,7 +560,7 @@ async fn responses_websocket_emits_websocket_telemetry_events() { tokio::time::sleep(Duration::from_millis(10)).await; let summary = harness - .otel_manager + .session_telemetry .runtime_metrics_summary() .expect("runtime metrics summary"); assert_eq!(summary.api_calls.count, 0); @@ -593,7 +593,7 @@ async fn responses_websocket_includes_timing_metrics_header_when_runtime_metrics .await; let harness = websocket_harness_with_runtime_metrics(&server, true).await; - harness.otel_manager.reset_runtime_metrics(); + harness.session_telemetry.reset_runtime_metrics(); let mut client_session = harness.client.new_session(); let prompt = prompt_with_input(vec![message_item("hello")]); @@ -607,7 +607,7 @@ async fn responses_websocket_includes_timing_metrics_header_when_runtime_metrics ); let summary = harness - .otel_manager + .session_telemetry .runtime_metrics_summary() .expect("runtime metrics summary"); assert_eq!(summary.responses_api_overhead_ms, 120); @@ -664,7 +664,7 @@ async fn responses_websocket_emits_reasoning_included_event() { .stream( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -736,7 +736,7 @@ async fn responses_websocket_emits_rate_limit_events() { .stream( &prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -1316,7 +1316,7 @@ async fn responses_websocket_v2_after_error_uses_full_create_without_previous_re .stream( &prompt_two, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, None, @@ -1515,7 +1515,7 @@ async fn websocket_harness_with_options( .with_runtime_reader(), ) .expect("in-memory metrics client"); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( conversation_id, MODEL, model_info.slug.as_str(), @@ -1548,7 +1548,7 @@ async fn websocket_harness_with_options( model_info, effort, summary, - otel_manager, + session_telemetry, } } @@ -1581,7 +1581,7 @@ async fn stream_until_complete_with_turn_metadata( .stream( prompt, &harness.model_info, - &harness.otel_manager, + &harness.session_telemetry, harness.effort, harness.summary, service_tier, diff --git a/codex-rs/otel/README.md b/codex-rs/otel/README.md index 02d4fc3a0..be90d6141 100644 --- a/codex-rs/otel/README.md +++ b/codex-rs/otel/README.md @@ -4,7 +4,7 @@ - Provider wiring for log/trace/metric exporters (`codex_otel::OtelProvider`, `codex_otel::provider`, and the compatibility shim `codex_otel::otel_provider`). -- Session-scoped business event emission via `codex_otel::OtelManager`. +- Session-scoped business event emission via `codex_otel::SessionTelemetry`. - Low-level metrics APIs via `codex_otel::metrics`. - Trace-context helpers via `codex_otel::trace_context` and crate-root re-exports. @@ -49,16 +49,16 @@ if let Some(provider) = OtelProvider::from(&settings)? { } ``` -## OtelManager (events) +## SessionTelemetry (events) -`OtelManager` adds consistent metadata to tracing events and helps record +`SessionTelemetry` adds consistent metadata to tracing events and helps record Codex-specific session events. Rich session/business events should go through -`OtelManager`; subsystem-owned audit events can stay with the owning subsystem. +`SessionTelemetry`; subsystem-owned audit events can stay with the owning subsystem. ```rust -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; -let manager = OtelManager::new( +let manager = SessionTelemetry::new( conversation_id, model, slug, @@ -134,7 +134,7 @@ use codex_otel::set_parent_from_w3c_trace_context; ## Shutdown - `OtelProvider::shutdown()` stops the OTEL exporter. -- `OtelManager::shutdown_metrics()` flushes and shuts down the metrics provider. +- `SessionTelemetry::shutdown_metrics()` flushes and shuts down the metrics provider. Both are optional because drop performs best-effort shutdown, but calling them explicitly gives deterministic flushing (or a shutdown error if flushing does diff --git a/codex-rs/otel/src/events/mod.rs b/codex-rs/otel/src/events/mod.rs index cd2c922e5..b0254c92c 100644 --- a/codex-rs/otel/src/events/mod.rs +++ b/codex-rs/otel/src/events/mod.rs @@ -1,2 +1,2 @@ -pub(crate) mod otel_manager; +pub(crate) mod session_telemetry; pub(crate) mod shared; diff --git a/codex-rs/otel/src/events/otel_manager.rs b/codex-rs/otel/src/events/session_telemetry.rs similarity index 98% rename from codex-rs/otel/src/events/otel_manager.rs rename to codex-rs/otel/src/events/session_telemetry.rs index bfe477f56..50d7614a0 100644 --- a/codex-rs/otel/src/events/otel_manager.rs +++ b/codex-rs/otel/src/events/session_telemetry.rs @@ -64,7 +64,7 @@ const RESPONSES_API_ENGINE_IAPI_TBT_FIELD: &str = "engine_iapi_tbt_across_engine const RESPONSES_API_ENGINE_SERVICE_TBT_FIELD: &str = "engine_service_tbt_across_engine_calls_ms"; #[derive(Debug, Clone)] -pub struct OtelEventMetadata { +pub struct SessionTelemetryMetadata { pub(crate) conversation_id: ThreadId, pub(crate) auth_mode: Option, pub(crate) account_id: Option, @@ -80,13 +80,13 @@ pub struct OtelEventMetadata { } #[derive(Debug, Clone)] -pub struct OtelManager { - pub(crate) metadata: OtelEventMetadata, +pub struct SessionTelemetry { + pub(crate) metadata: SessionTelemetryMetadata, pub(crate) metrics: Option, pub(crate) metrics_use_metadata_tags: bool, } -impl OtelManager { +impl SessionTelemetry { pub fn with_model(mut self, model: &str, slug: &str) -> Self { self.metadata.model = model.to_owned(); self.metadata.slug = slug.to_owned(); @@ -276,9 +276,9 @@ impl OtelManager { log_user_prompts: bool, terminal_type: String, session_source: SessionSource, - ) -> OtelManager { + ) -> SessionTelemetry { Self { - metadata: OtelEventMetadata { + metadata: SessionTelemetryMetadata { conversation_id, auth_mode: auth_mode.map(|m| m.to_string()), account_id, @@ -298,7 +298,7 @@ impl OtelManager { } pub fn record_responses(&self, handle_responses_span: &Span, event: &ResponseEvent) { - handle_responses_span.record("otel.name", OtelManager::responses_type(event)); + handle_responses_span.record("otel.name", SessionTelemetry::responses_type(event)); match event { ResponseEvent::OutputItemDone(item) => { @@ -902,7 +902,7 @@ impl OtelManager { match event { ResponseEvent::Created => "created".into(), ResponseEvent::OutputItemDone(item) | ResponseEvent::OutputItemAdded(item) => { - OtelManager::responses_item_type(item) + SessionTelemetry::responses_item_type(item) } ResponseEvent::Completed { .. } => "completed".into(), ResponseEvent::OutputTextDelta(_) => "text_delta".into(), diff --git a/codex-rs/otel/src/lib.rs b/codex-rs/otel/src/lib.rs index 5594b26a7..5a4ba31e4 100644 --- a/codex-rs/otel/src/lib.rs +++ b/codex-rs/otel/src/lib.rs @@ -13,8 +13,8 @@ use crate::metrics::Result as MetricsResult; use serde::Serialize; use strum_macros::Display; -pub use crate::events::otel_manager::OtelEventMetadata; -pub use crate::events::otel_manager::OtelManager; +pub use crate::events::session_telemetry::SessionTelemetry; +pub use crate::events::session_telemetry::SessionTelemetryMetadata; 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/tests/suite/manager_metrics.rs b/codex-rs/otel/tests/suite/manager_metrics.rs index 53a9cc89d..a3160d5a4 100644 --- a/codex-rs/otel/tests/suite/manager_metrics.rs +++ b/codex-rs/otel/tests/suite/manager_metrics.rs @@ -2,7 +2,7 @@ use crate::harness::attributes_to_map; use crate::harness::build_metrics_with_defaults; use crate::harness::find_metric; use crate::harness::latest_metrics; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_otel::metrics::Result; use codex_protocol::ThreadId; @@ -12,11 +12,11 @@ use opentelemetry_sdk::metrics::data::MetricData; use pretty_assertions::assert_eq; use std::collections::BTreeMap; -// Ensures OtelManager attaches metadata tags when forwarding metrics. +// Ensures SessionTelemetry attaches metadata tags when forwarding metrics. #[test] fn manager_attaches_metadata_tags_to_metrics() -> Result<()> { let (metrics, exporter) = build_metrics_with_defaults(&[("service", "codex-cli")])?; - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", @@ -68,11 +68,11 @@ fn manager_attaches_metadata_tags_to_metrics() -> Result<()> { Ok(()) } -// Ensures metadata tagging can be disabled when recording via OtelManager. +// Ensures metadata tagging can be disabled when recording via SessionTelemetry. #[test] fn manager_allows_disabling_metadata_tags() -> Result<()> { let (metrics, exporter) = build_metrics_with_defaults(&[])?; - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-4o", "gpt-4o", @@ -113,7 +113,7 @@ fn manager_allows_disabling_metadata_tags() -> Result<()> { #[test] fn manager_attaches_optional_service_name_tag() -> Result<()> { let (metrics, exporter) = build_metrics_with_defaults(&[])?; - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", diff --git a/codex-rs/otel/tests/suite/otel_export_routing_policy.rs b/codex-rs/otel/tests/suite/otel_export_routing_policy.rs index 875bfd666..75d9bde83 100644 --- a/codex-rs/otel/tests/suite/otel_export_routing_policy.rs +++ b/codex-rs/otel/tests/suite/otel_export_routing_policy.rs @@ -1,4 +1,4 @@ -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_otel::otel_provider::OtelProvider; use opentelemetry::KeyValue; @@ -103,7 +103,7 @@ fn otel_export_routing_policy_routes_user_prompt_log_and_trace_events() { tracing::subscriber::with_default(subscriber, || { tracing::callsite::rebuild_interest_cache(); - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", @@ -212,7 +212,7 @@ fn otel_export_routing_policy_routes_tool_result_log_and_trace_events() { tracing::subscriber::with_default(subscriber, || { tracing::callsite::rebuild_interest_cache(); - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", diff --git a/codex-rs/otel/tests/suite/runtime_summary.rs b/codex-rs/otel/tests/suite/runtime_summary.rs index 8857387d5..c2f252381 100644 --- a/codex-rs/otel/tests/suite/runtime_summary.rs +++ b/codex-rs/otel/tests/suite/runtime_summary.rs @@ -1,6 +1,6 @@ -use codex_otel::OtelManager; use codex_otel::RuntimeMetricTotals; use codex_otel::RuntimeMetricsSummary; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_otel::metrics::MetricsClient; use codex_otel::metrics::MetricsConfig; @@ -20,7 +20,7 @@ fn runtime_metrics_summary_collects_tool_api_and_streaming_metrics() -> Result<( MetricsConfig::in_memory("test", "codex-cli", env!("CARGO_PKG_VERSION"), exporter) .with_runtime_reader(), )?; - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", diff --git a/codex-rs/otel/tests/suite/snapshot.rs b/codex-rs/otel/tests/suite/snapshot.rs index e13775af6..6bd9bd819 100644 --- a/codex-rs/otel/tests/suite/snapshot.rs +++ b/codex-rs/otel/tests/suite/snapshot.rs @@ -1,6 +1,6 @@ use crate::harness::attributes_to_map; use crate::harness::find_metric; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_otel::metrics::MetricsClient; use codex_otel::metrics::MetricsConfig; @@ -69,7 +69,7 @@ fn manager_snapshot_metrics_collects_without_shutdown() -> Result<()> { .with_tag("service", "codex-cli")? .with_runtime_reader(); let metrics = MetricsClient::new(config)?; - let manager = OtelManager::new( + let manager = SessionTelemetry::new( ThreadId::new(), "gpt-5.1", "gpt-5.1", diff --git a/codex-rs/state/src/runtime.rs b/codex-rs/state/src/runtime.rs index a32f041ce..647565b73 100644 --- a/codex-rs/state/src/runtime.rs +++ b/codex-rs/state/src/runtime.rs @@ -28,7 +28,7 @@ use crate::model::datetime_to_epoch_seconds; use crate::paths::file_modified_time_utc; use chrono::DateTime; use chrono::Utc; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::dynamic_tools::DynamicToolSpec; use codex_protocol::protocol::RolloutItem; @@ -84,7 +84,7 @@ impl StateRuntime { pub async fn init( codex_home: PathBuf, default_provider: String, - otel: Option, + otel: Option, ) -> anyhow::Result> { tokio::fs::create_dir_all(&codex_home).await?; let current_state_name = state_db_filename(); diff --git a/codex-rs/state/src/runtime/threads.rs b/codex-rs/state/src/runtime/threads.rs index e671dfaf6..8825f13fe 100644 --- a/codex-rs/state/src/runtime/threads.rs +++ b/codex-rs/state/src/runtime/threads.rs @@ -447,7 +447,7 @@ ON CONFLICT(thread_id, position) DO NOTHING &self, builder: &ThreadMetadataBuilder, items: &[RolloutItem], - otel: Option<&OtelManager>, + otel: Option<&SessionTelemetry>, new_thread_memory_mode: Option<&str>, updated_at_override: Option>, ) -> anyhow::Result<()> { diff --git a/codex-rs/tui/src/app.rs b/codex-rs/tui/src/app.rs index 32f5d5963..9608cf563 100644 --- a/codex-rs/tui/src/app.rs +++ b/codex-rs/tui/src/app.rs @@ -57,7 +57,7 @@ use codex_core::models_manager::model_presets::HIDE_GPT_5_1_CODEX_MAX_MIGRATION_ use codex_core::models_manager::model_presets::HIDE_GPT5_1_MIGRATION_PROMPT_CONFIG; #[cfg(target_os = "windows")] use codex_core::windows_sandbox::WindowsSandboxLevelExt; -use codex_otel::OtelManager; +use codex_otel::SessionTelemetry; use codex_otel::TelemetryAuthMode; use codex_protocol::ThreadId; use codex_protocol::config_types::Personality; @@ -637,7 +637,7 @@ async fn handle_model_migration_prompt_if_needed( pub(crate) struct App { pub(crate) server: Arc, - pub(crate) otel_manager: OtelManager, + pub(crate) session_telemetry: SessionTelemetry, pub(crate) app_event_tx: AppEventSender, pub(crate) chat_widget: ChatWidget, pub(crate) auth_manager: Arc, @@ -750,7 +750,7 @@ impl App { model: Some(self.chat_widget.current_model().to_string()), startup_tooltip_override: None, status_line_invalid_items_warned: self.status_line_invalid_items_warned.clone(), - otel_manager: self.otel_manager.clone(), + session_telemetry: self.session_telemetry.clone(), } } @@ -1437,7 +1437,7 @@ impl App { model: Some(model), startup_tooltip_override: None, status_line_invalid_items_warned: self.status_line_invalid_items_warned.clone(), - otel_manager: self.otel_manager.clone(), + session_telemetry: self.session_telemetry.clone(), }; self.chat_widget = ChatWidget::new(init, self.server.clone()); self.reset_thread_event_state(); @@ -1624,7 +1624,7 @@ impl App { let auth_mode = auth_ref .map(CodexAuth::auth_mode) .map(TelemetryAuthMode::from); - let otel_manager = OtelManager::new( + let session_telemetry = SessionTelemetry::new( ThreadId::new(), model.as_str(), model.as_str(), @@ -1641,7 +1641,7 @@ impl App { .as_ref() .is_some_and(|cmd| !cmd.is_empty()) { - otel_manager.counter("codex.status_line", 1, &[]); + session_telemetry.counter("codex.status_line", 1, &[]); } let status_line_invalid_items_warned = Arc::new(AtomicBool::new(false)); @@ -1673,7 +1673,7 @@ impl App { model: Some(model.clone()), startup_tooltip_override, status_line_invalid_items_warned: status_line_invalid_items_warned.clone(), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), }; ChatWidget::new(init, thread_manager.clone()) } @@ -1708,12 +1708,12 @@ impl App { model: config.model.clone(), startup_tooltip_override: None, status_line_invalid_items_warned: status_line_invalid_items_warned.clone(), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), }; ChatWidget::new_from_existing(init, resumed.thread, resumed.session_configured) } SessionSelection::Fork(target_session) => { - otel_manager.counter("codex.thread.fork", 1, &[("source", "cli_subcommand")]); + session_telemetry.counter("codex.thread.fork", 1, &[("source", "cli_subcommand")]); let forked = thread_manager .fork_thread( usize::MAX, @@ -1745,7 +1745,7 @@ impl App { model: config.model.clone(), startup_tooltip_override: None, status_line_invalid_items_warned: status_line_invalid_items_warned.clone(), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), }; ChatWidget::new_from_existing(init, forked.thread, forked.session_configured) } @@ -1760,7 +1760,7 @@ impl App { let mut app = Self { server: thread_manager.clone(), - otel_manager: otel_manager.clone(), + session_telemetry: session_telemetry.clone(), app_event_tx, chat_widget, auth_manager: auth_manager.clone(), @@ -2081,8 +2081,11 @@ impl App { tui.frame_requester().schedule_frame(); } AppEvent::ForkCurrentSession => { - self.otel_manager - .counter("codex.thread.fork", 1, &[("source", "slash_command")]); + self.session_telemetry.counter( + "codex.thread.fork", + 1, + &[("source", "slash_command")], + ); let summary = session_summary( self.chat_widget.token_usage(), self.chat_widget.thread_id(), @@ -2343,11 +2346,14 @@ impl App { self.chat_widget.open_windows_sandbox_enable_prompt(preset); } AppEvent::OpenWindowsSandboxFallbackPrompt { preset } => { - self.otel_manager - .counter("codex.windows_sandbox.fallback_prompt_shown", 1, &[]); + self.session_telemetry.counter( + "codex.windows_sandbox.fallback_prompt_shown", + 1, + &[], + ); self.chat_widget.clear_windows_sandbox_setup_status(); if let Some(started_at) = self.windows_sandbox.setup_started_at.take() { - self.otel_manager.record_duration( + self.session_telemetry.record_duration( "codex.windows_sandbox.elevated_setup_duration_ms", started_at.elapsed(), &[("result", "failure")], @@ -2380,7 +2386,7 @@ impl App { self.chat_widget.show_windows_sandbox_setup_status(); self.windows_sandbox.setup_started_at = Some(Instant::now()); - let otel_manager = self.otel_manager.clone(); + let session_telemetry = self.session_telemetry.clone(); tokio::task::spawn_blocking(move || { let result = codex_core::windows_sandbox::run_elevated_setup( &policy, @@ -2391,7 +2397,7 @@ impl App { ); let event = match result { Ok(()) => { - otel_manager.counter( + session_telemetry.counter( "codex.windows_sandbox.elevated_setup_success", 1, &[], @@ -2419,7 +2425,7 @@ impl App { if let Some(message) = message_tag.as_deref() { tags.push(("message", message)); } - otel_manager.counter( + session_telemetry.counter( codex_core::windows_sandbox::elevated_setup_failure_metric_name( &err, ), @@ -2451,7 +2457,7 @@ impl App { std::env::vars().collect(); let codex_home = self.config.codex_home.clone(); let tx = self.app_event_tx.clone(); - let otel_manager = self.otel_manager.clone(); + let session_telemetry = self.session_telemetry.clone(); self.chat_widget.show_windows_sandbox_setup_status(); tokio::task::spawn_blocking(move || { @@ -2462,7 +2468,7 @@ impl App { &env_map, codex_home.as_path(), ) { - otel_manager.counter( + session_telemetry.counter( "codex.windows_sandbox.legacy_setup_preflight_failed", 1, &[], @@ -2545,7 +2551,7 @@ impl App { { self.chat_widget.clear_windows_sandbox_setup_status(); if let Some(started_at) = self.windows_sandbox.setup_started_at.take() { - self.otel_manager.record_duration( + self.session_telemetry.record_duration( "codex.windows_sandbox.elevated_setup_duration_ms", started_at.elapsed(), &[("result", "success")], @@ -3713,7 +3719,7 @@ mod tests { use codex_core::config::ConfigBuilder; use codex_core::config::ConfigOverrides; use codex_core::config::types::ModelAvailabilityNuxConfig; - use codex_otel::OtelManager; + use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::config_types::CollaborationMode; use codex_protocol::config_types::CollaborationModeMask; @@ -5276,11 +5282,11 @@ mod tests { ); let file_search = FileSearchManager::new(config.cwd.clone(), app_event_tx.clone()); let model = codex_core::test_support::get_model_offline(config.model.as_deref()); - let otel_manager = test_otel_manager(&config, model.as_str()); + let session_telemetry = test_session_telemetry(&config, model.as_str()); App { server, - otel_manager, + session_telemetry, app_event_tx, chat_widget, auth_manager, @@ -5335,12 +5341,12 @@ mod tests { ); let file_search = FileSearchManager::new(config.cwd.clone(), app_event_tx.clone()); let model = codex_core::test_support::get_model_offline(config.model.as_deref()); - let otel_manager = test_otel_manager(&config, model.as_str()); + let session_telemetry = test_session_telemetry(&config, model.as_str()); ( App { server, - otel_manager, + session_telemetry, app_event_tx, chat_widget, auth_manager, @@ -5391,9 +5397,9 @@ mod tests { panic!("expected UserTurn op, saw: {seen:?}"); } - fn test_otel_manager(config: &Config, model: &str) -> OtelManager { + fn test_session_telemetry(config: &Config, model: &str) -> SessionTelemetry { let model_info = codex_core::test_support::construct_model_info_offline(model, config); - OtelManager::new( + SessionTelemetry::new( ThreadId::new(), model, model_info.slug.as_str(), diff --git a/codex-rs/tui/src/chatwidget.rs b/codex-rs/tui/src/chatwidget.rs index 9f2295e65..ae35225b6 100644 --- a/codex-rs/tui/src/chatwidget.rs +++ b/codex-rs/tui/src/chatwidget.rs @@ -74,8 +74,8 @@ use codex_core::terminal::TerminalName; use codex_core::terminal::terminal_info; #[cfg(target_os = "windows")] use codex_core::windows_sandbox::WindowsSandboxLevelExt; -use codex_otel::OtelManager; use codex_otel::RuntimeMetricsSummary; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::account::PlanType; use codex_protocol::approvals::ElicitationRequestEvent; @@ -478,7 +478,7 @@ pub(crate) struct ChatWidgetInit { pub(crate) startup_tooltip_override: Option, // Shared latch so we only warn once about invalid status-line item IDs. pub(crate) status_line_invalid_items_warned: Arc, - pub(crate) otel_manager: OtelManager, + pub(crate) session_telemetry: SessionTelemetry, } #[derive(Default)] @@ -560,7 +560,7 @@ pub(crate) struct ChatWidget { active_collaboration_mask: Option, auth_manager: Arc, models_manager: Arc, - otel_manager: OtelManager, + session_telemetry: SessionTelemetry, session_header: SessionHeader, initial_user_message: Option, token_info: Option, @@ -1145,7 +1145,7 @@ impl ChatWidget { } fn collect_runtime_metrics_delta(&mut self) { - if let Some(delta) = self.otel_manager.runtime_metrics_summary() { + if let Some(delta) = self.session_telemetry.runtime_metrics_summary() { self.apply_runtime_metrics_delta(delta); } } @@ -1506,7 +1506,7 @@ impl ChatWidget { self.adaptive_chunking.reset(); self.plan_stream_controller = None; self.turn_runtime_metrics = RuntimeMetricsSummary::default(); - self.otel_manager.reset_runtime_metrics(); + self.session_telemetry.reset_runtime_metrics(); self.bottom_pane.clear_quit_shortcut_hint(); self.quit_shortcut_expires_at = None; self.quit_shortcut_key = None; @@ -3044,7 +3044,7 @@ impl ChatWidget { model, startup_tooltip_override, status_line_invalid_items_warned, - otel_manager, + session_telemetry, } = common; let model = model.filter(|m| !m.trim().is_empty()); let mut config = config; @@ -3103,7 +3103,7 @@ impl ChatWidget { active_collaboration_mask, auth_manager, models_manager, - otel_manager, + session_telemetry, session_header: SessionHeader::new(header_model), initial_user_message, token_info: None, @@ -3227,7 +3227,7 @@ impl ChatWidget { model, startup_tooltip_override, status_line_invalid_items_warned, - otel_manager, + session_telemetry, } = common; let model = model.filter(|m| !m.trim().is_empty()); let mut config = config; @@ -3285,7 +3285,7 @@ impl ChatWidget { active_collaboration_mask, auth_manager, models_manager, - otel_manager, + session_telemetry, session_header: SessionHeader::new(header_model), initial_user_message, token_info: None, @@ -3401,7 +3401,7 @@ impl ChatWidget { model, startup_tooltip_override: _, status_line_invalid_items_warned, - otel_manager, + session_telemetry, } = common; let model = model.filter(|m| !m.trim().is_empty()); let prevent_idle_sleep = config.features.enabled(Feature::PreventIdleSleep); @@ -3459,7 +3459,7 @@ impl ChatWidget { active_collaboration_mask, auth_manager, models_manager, - otel_manager, + session_telemetry, session_header: SessionHeader::new(header_model), initial_user_message, token_info: None, @@ -3840,7 +3840,8 @@ impl ChatWidget { self.open_review_popup(); } SlashCommand::Rename => { - self.otel_manager.counter("codex.thread.rename", 1, &[]); + self.session_telemetry + .counter("codex.thread.rename", 1, &[]); self.show_rename_prompt(); } SlashCommand::Model => { @@ -3942,7 +3943,7 @@ impl ChatWidget { return; } - self.otel_manager.counter( + self.session_telemetry.counter( "codex.windows_sandbox.setup_elevated_sandbox_command", 1, &[], @@ -3952,7 +3953,7 @@ impl ChatWidget { } #[cfg(not(target_os = "windows"))] { - let _ = &self.otel_manager; + let _ = &self.session_telemetry; // Not supported; on non-Windows this command should never be reachable. }; } @@ -4156,7 +4157,8 @@ impl ChatWidget { } } SlashCommand::Rename if !trimmed.is_empty() => { - self.otel_manager.counter("codex.thread.rename", 1, &[]); + self.session_telemetry + .counter("codex.thread.rename", 1, &[]); let Some((prepared_args, _prepared_elements)) = self.bottom_pane.prepare_inline_args_submission(false) else { @@ -6929,7 +6931,7 @@ impl ChatWidget { return; } - self.otel_manager + self.session_telemetry .counter("codex.windows_sandbox.elevated_prompt_shown", 1, &[]); let mut header = ColumnRenderable::new(); @@ -6940,10 +6942,10 @@ impl ChatWidget { .wrap(Wrap { trim: false }), )); - let accept_otel = self.otel_manager.clone(); - let legacy_otel = self.otel_manager.clone(); + let accept_otel = self.session_telemetry.clone(); + let legacy_otel = self.session_telemetry.clone(); let legacy_preset = preset.clone(); - let quit_otel = self.otel_manager.clone(); + let quit_otel = self.session_telemetry.clone(); let items = vec![ SelectionItem { name: "Set up default sandbox (requires Administrator permissions)".to_string(), @@ -7014,13 +7016,13 @@ impl ChatWidget { let elevated_preset = preset.clone(); let legacy_preset = preset; - let quit_otel = self.otel_manager.clone(); + let quit_otel = self.session_telemetry.clone(); let items = vec![ SelectionItem { name: "Try setting up admin sandbox again".to_string(), description: None, actions: vec![Box::new({ - let otel = self.otel_manager.clone(); + let otel = self.session_telemetry.clone(); let preset = elevated_preset; move |tx| { otel.counter("codex.windows_sandbox.fallback_retry_elevated", 1, &[]); @@ -7036,7 +7038,7 @@ impl ChatWidget { name: "Use Codex with non-admin sandbox".to_string(), description: None, actions: vec![Box::new({ - let otel = self.otel_manager.clone(); + let otel = self.session_telemetry.clone(); let preset = legacy_preset; move |tx| { otel.counter("codex.windows_sandbox.fallback_use_legacy", 1, &[]); diff --git a/codex-rs/tui/src/chatwidget/tests.rs b/codex-rs/tui/src/chatwidget/tests.rs index ddcd015c6..a9819fb18 100644 --- a/codex-rs/tui/src/chatwidget/tests.rs +++ b/codex-rs/tui/src/chatwidget/tests.rs @@ -32,8 +32,8 @@ use codex_core::models_manager::collaboration_mode_presets::CollaborationModesCo use codex_core::models_manager::manager::ModelsManager; use codex_core::skills::model::SkillMetadata; use codex_core::terminal::TerminalName; -use codex_otel::OtelManager; use codex_otel::RuntimeMetricsSummary; +use codex_otel::SessionTelemetry; use codex_protocol::ThreadId; use codex_protocol::account::PlanType; use codex_protocol::config_types::CollaborationMode; @@ -1669,7 +1669,7 @@ async fn helpers_are_available_and_do_not_panic() { let tx = AppEventSender::new(tx_raw); let cfg = test_config().await; let resolved_model = codex_core::test_support::get_model_offline(cfg.model.as_deref()); - let otel_manager = test_otel_manager(&cfg, resolved_model.as_str()); + let session_telemetry = test_session_telemetry(&cfg, resolved_model.as_str()); let thread_manager = Arc::new( codex_core::test_support::thread_manager_with_models_provider( CodexAuth::from_api_key("test"), @@ -1692,16 +1692,16 @@ async fn helpers_are_available_and_do_not_panic() { model: Some(resolved_model), startup_tooltip_override: None, status_line_invalid_items_warned: Arc::new(AtomicBool::new(false)), - otel_manager, + session_telemetry, }; let mut w = ChatWidget::new(init, thread_manager); // Basic construction sanity. let _ = &mut w; } -fn test_otel_manager(config: &Config, model: &str) -> OtelManager { +fn test_session_telemetry(config: &Config, model: &str) -> SessionTelemetry { let model_info = codex_core::test_support::construct_model_info_offline(model, config); - OtelManager::new( + SessionTelemetry::new( ThreadId::new(), model, model_info.slug.as_str(), @@ -1734,7 +1734,7 @@ async fn make_chatwidget_manual( cfg.model = Some(model.to_string()); } let prevent_idle_sleep = cfg.features.enabled(Feature::PreventIdleSleep); - let otel_manager = test_otel_manager(&cfg, resolved_model.as_str()); + let session_telemetry = test_session_telemetry(&cfg, resolved_model.as_str()); let mut bottom = BottomPane::new(BottomPaneParams { app_event_tx: app_event_tx.clone(), frame_requester: FrameRequester::test_dummy(), @@ -1777,7 +1777,7 @@ async fn make_chatwidget_manual( active_collaboration_mask, auth_manager, models_manager, - otel_manager, + session_telemetry, session_header: SessionHeader::new(resolved_model.clone()), initial_user_message: None, token_info: None, @@ -5359,7 +5359,7 @@ async fn collaboration_modes_defaults_to_code_on_startup() { .await .expect("config"); let resolved_model = codex_core::test_support::get_model_offline(cfg.model.as_deref()); - let otel_manager = test_otel_manager(&cfg, resolved_model.as_str()); + let session_telemetry = test_session_telemetry(&cfg, resolved_model.as_str()); let thread_manager = Arc::new( codex_core::test_support::thread_manager_with_models_provider( CodexAuth::from_api_key("test"), @@ -5382,7 +5382,7 @@ async fn collaboration_modes_defaults_to_code_on_startup() { model: Some(resolved_model.clone()), startup_tooltip_override: None, status_line_invalid_items_warned: Arc::new(AtomicBool::new(false)), - otel_manager, + session_telemetry, }; let chat = ChatWidget::new(init, thread_manager); @@ -5409,7 +5409,7 @@ async fn experimental_mode_plan_is_ignored_on_startup() { .await .expect("config"); let resolved_model = codex_core::test_support::get_model_offline(cfg.model.as_deref()); - let otel_manager = test_otel_manager(&cfg, resolved_model.as_str()); + let session_telemetry = test_session_telemetry(&cfg, resolved_model.as_str()); let thread_manager = Arc::new( codex_core::test_support::thread_manager_with_models_provider( CodexAuth::from_api_key("test"), @@ -5432,7 +5432,7 @@ async fn experimental_mode_plan_is_ignored_on_startup() { model: Some(resolved_model.clone()), startup_tooltip_override: None, status_line_invalid_items_warned: Arc::new(AtomicBool::new(false)), - otel_manager, + session_telemetry, }; let chat = ChatWidget::new(init, thread_manager);