chore(otel): rename OtelManager to SessionTelemetry (#13808)

## Summary
This is a purely mechanical refactor of `OtelManager` ->
`SessionTelemetry` to better convey what the struct is doing. No
behavior change.

## Why

`OtelManager` ended up sounding much broader than what this type
actually does. It doesn't manage OTEL globally; it's the session-scoped
telemetry surface for emitting log/trace events and recording metrics
with consistent session metadata (`app_version`, `model`, `slug`,
`originator`, etc.).

`SessionTelemetry` is a more accurate name, and updating the call sites
makes that boundary a lot easier to follow.

## Validation

- `just fmt`
- `cargo test -p codex-otel`
- `cargo test -p codex-core`
This commit is contained in:
Owen Lin
2026-03-06 16:23:30 -08:00
committed by GitHub
Unverified
parent 3794363cac
commit 289ed549cf
45 changed files with 318 additions and 290 deletions
+52 -48
View File
@@ -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<Vec<ResponseItem>> {
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<ApiRawMemory>,
model_info: &ModelInfo,
effort: Option<ReasoningEffortConfig>,
otel_manager: &OtelManager,
session_telemetry: &SessionTelemetry,
) -> Result<Vec<ApiMemorySummarizeOutput>> {
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<dyn RequestTelemetry> {
let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone()));
fn build_request_telemetry(session_telemetry: &SessionTelemetry) -> Arc<dyn RequestTelemetry> {
let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone()));
let request_telemetry: Arc<dyn RequestTelemetry> = 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<Arc<OnceLock<String>>>,
turn_metadata_header: Option<&str>,
) -> std::result::Result<ApiWebSocketConnection, ApiError> {
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<ReasoningEffortConfig>,
summary: ReasoningSummaryConfig,
service_tier: Option<ServiceTier>,
@@ -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<ReasoningEffortConfig>,
summary: ReasoningSummaryConfig,
service_tier: Option<ServiceTier>,
@@ -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<dyn RequestTelemetry>, Arc<dyn SseTelemetry>) {
let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone()));
let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone()));
let request_telemetry: Arc<dyn RequestTelemetry> = telemetry.clone();
let sse_telemetry: Arc<dyn SseTelemetry> = telemetry;
(request_telemetry, sse_telemetry)
}
/// Builds telemetry for the Responses API WebSocket transport.
fn build_websocket_telemetry(otel_manager: &OtelManager) -> Arc<dyn WebsocketTelemetry> {
let telemetry = Arc::new(ApiTelemetry::new(otel_manager.clone()));
fn build_websocket_telemetry(
session_telemetry: &SessionTelemetry,
) -> Arc<dyn WebsocketTelemetry> {
let telemetry = Arc::new(ApiTelemetry::new(session_telemetry.clone()));
let websocket_telemetry: Arc<dyn WebsocketTelemetry> = telemetry;
websocket_telemetry
}
@@ -918,7 +921,7 @@ impl ModelClientSession {
&mut self,
prompt: &Prompt,
model_info: &ModelInfo,
otel_manager: &OtelManager,
session_telemetry: &SessionTelemetry,
effort: Option<ReasoningEffortConfig>,
summary: ReasoningSummaryConfig,
service_tier: Option<ServiceTier>,
@@ -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<ReasoningEffortConfig>,
summary: ReasoningSummaryConfig,
service_tier: Option<ServiceTier>,
@@ -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<S>(
api_stream: S,
otel_manager: OtelManager,
session_telemetry: SessionTelemetry,
) -> (ResponseStream, oneshot::Receiver<LastResponse>)
where
S: futures::Stream<Item = std::result::Result<ResponseEvent, ApiError>>
@@ -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<Option<std::result::Result<Message, Error>>, 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);
+32 -30
View File
@@ -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<Config>,
pub(crate) auth_manager: Option<Arc<AuthManager>>,
pub(crate) model_info: ModelInfo,
pub(crate) otel_manager: OtelManager,
pub(crate) session_telemetry: SessionTelemetry,
pub(crate) provider: ModelProviderInfo,
pub(crate) reasoning_effort: Option<ReasoningEffortConfig>,
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<Arc<AuthManager>>,
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<String>, 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;
+9 -9
View File
@@ -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,
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -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;
+2 -2
View File
@@ -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;
+3 -3
View File
@@ -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 }
+17 -17
View File
@@ -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<ReasoningEffortConfig>,
pub(in crate::memories) reasoning_summary: ReasoningSummaryConfig,
pub(in crate::memories) service_tier: Option<ServiceTier>,
@@ -85,7 +85,7 @@ struct StageOneOutput {
pub(in crate::memories) async fn run(session: &Arc<Session>, 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<Session>, 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<JobResult>) -> 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")],
+12 -8
View File
@@ -43,7 +43,7 @@ struct Counters {
pub(super) async fn run(session: &Arc<Session>, config: Arc<Config>) {
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<Session>, config: Arc<Config>) {
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<Session>,
db: &StateRuntime,
) -> Result<Claim, &'static str> {
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<Session>, 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<Session>, counters: Counters) {
}
fn emit_token_usage_metrics(session: &Arc<Session>, 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),
+1 -1
View File
@@ -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,
&[
+3 -3
View File
@@ -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<ReasoningEffortConfig>,
otel_manager: &OtelManager,
session_telemetry: &SessionTelemetry,
) -> Result<Vec<BuiltMemory>> {
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!(
+3 -3
View File
@@ -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<ExtractionOutcome> {
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 {
+8 -8
View File
@@ -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<Option<Arc<ShellSnapshot>>> {
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<Option<Arc<ShellSnapshot>>>,
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<Option<Arc<ShellSnapshot>>>,
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),
+7 -3
View File
@@ -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;
};
+1 -1
View File
@@ -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,
&[
+2 -2
View File
@@ -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<AuthManager>,
pub(crate) models_manager: Arc<ModelsManager>,
pub(crate) otel_manager: OtelManager,
pub(crate) session_telemetry: SessionTelemetry,
pub(crate) tool_approvals: Mutex<ApprovalStore>,
#[cfg_attr(not(unix), allow(dead_code))]
pub(crate) execve_session_approvals: RwLock<HashMap<AbsolutePathBuf, ExecveSessionApproval>>,
+9 -3
View File
@@ -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<codex_state::StateRuntime>;
/// 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<StateDbHandle> {
pub(crate) async fn init(
config: &Config,
otel: Option<&SessionTelemetry>,
) -> Option<StateDbHandle> {
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<StateDbHandle> {
pub async fn get_state_db(
config: &Config,
otel: Option<&SessionTelemetry>,
) -> Option<StateDbHandle> {
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;
+1 -1
View File
@@ -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);
+2 -2
View File
@@ -30,14 +30,14 @@ impl SessionTask for CompactTask {
) -> Option<String> {
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")],
+7 -7
View File
@@ -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],
+1 -1
View File
@@ -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,
+1 -1
View File
@@ -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.
+1 -1
View File
@@ -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(
+1 -1
View File
@@ -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 {
@@ -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| {
+1 -1
View File
@@ -108,7 +108,7 @@ impl ToolOrchestrator {
where
T: ToolRuntime<Rq, Out>,
{
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;
+1 -1
View File
@@ -82,7 +82,7 @@ impl ToolRegistry {
) -> Result<ResponseInputItem, FunctionCallError> {
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 = [
+1 -1
View File
@@ -88,7 +88,7 @@ where
let decision = fetch().await;
services.otel_manager.counter(
services.session_telemetry.counter(
"codex.approval.requested",
1,
&[
+2 -2
View File
@@ -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, &[]);
}
+7 -7
View File
@@ -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,
+3 -3
View File
@@ -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,
+20 -20
View File
@@ -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<ReasoningEffortConfig>,
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,
+7 -7
View File
@@ -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
+1 -1
View File
@@ -1,2 +1,2 @@
pub(crate) mod otel_manager;
pub(crate) mod session_telemetry;
pub(crate) mod shared;
@@ -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<String>,
pub(crate) account_id: Option<String>,
@@ -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<MetricsClient>,
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(),
+2 -2
View File
@@ -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;
+6 -6
View File
@@ -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",
@@ -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",
+2 -2
View File
@@ -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",
+2 -2
View File
@@ -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",
+2 -2
View File
@@ -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<OtelManager>,
otel: Option<SessionTelemetry>,
) -> anyhow::Result<Arc<Self>> {
tokio::fs::create_dir_all(&codex_home).await?;
let current_state_name = state_db_filename();
+1 -1
View File
@@ -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<DateTime<Utc>>,
) -> anyhow::Result<()> {
+35 -29
View File
@@ -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<ThreadManager>,
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<AuthManager>,
@@ -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(),
+24 -22
View File
@@ -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<String>,
// Shared latch so we only warn once about invalid status-line item IDs.
pub(crate) status_line_invalid_items_warned: Arc<AtomicBool>,
pub(crate) otel_manager: OtelManager,
pub(crate) session_telemetry: SessionTelemetry,
}
#[derive(Default)]
@@ -560,7 +560,7 @@ pub(crate) struct ChatWidget {
active_collaboration_mask: Option<CollaborationModeMask>,
auth_manager: Arc<AuthManager>,
models_manager: Arc<ModelsManager>,
otel_manager: OtelManager,
session_telemetry: SessionTelemetry,
session_header: SessionHeader,
initial_user_message: Option<UserMessage>,
token_info: Option<TokenUsageInfo>,
@@ -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, &[]);
+11 -11
View File
@@ -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);