diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index f6968e3101..b25d1f5c83 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -530,6 +530,7 @@ async fn ingest_rejected_turn_steer( response: Box::new(sample_thread_resume_response( "thread-2", /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, out, ) @@ -603,6 +604,7 @@ async fn ingest_turn_prerequisites( response: Box::new(sample_thread_start_response( "thread-2", /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, out, ) @@ -626,6 +628,7 @@ async fn ingest_turn_prerequisites( connection_id: 7, request_id: RequestId::Integer(3), response: Box::new(sample_turn_start_response("turn-2")), + thread_initialization: None, }, out, ) @@ -692,6 +695,7 @@ async fn ingest_review_prerequisites( response: Box::new(sample_thread_start_response( "thread-1", /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, events, ) @@ -1604,6 +1608,7 @@ async fn initialize_caches_client_and_thread_lifecycle_publishes_once_initialize /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, &mut events, ) @@ -1648,6 +1653,7 @@ async fn initialize_caches_client_and_thread_lifecycle_publishes_once_initialize response: Box::new(sample_thread_resume_response( "thread-1", /*ephemeral*/ true, "gpt-5", )), + thread_initialization: None, }, &mut events, ) @@ -1718,6 +1724,7 @@ async fn unrelated_client_requests_are_ignored_by_reducer() { connection_id: 7, request_id: RequestId::Integer(3), response: Box::new(sample_turn_start_response("turn-2")), + thread_initialization: None, }, &mut events, ) @@ -1743,6 +1750,7 @@ async fn unrelated_client_responses_are_ignored_by_reducer() { response: Box::new(ClientResponsePayload::ThreadArchive( ThreadArchiveResponse {}, )), + thread_initialization: None, }, &mut events, ) @@ -1801,6 +1809,7 @@ async fn compaction_event_ingests_custom_fact() { Some(AppServerThreadSource::Subagent), Some(parent_thread_id.to_string()), )), + thread_initialization: None, }, &mut events, ) @@ -1908,6 +1917,7 @@ async fn guardian_review_event_ingests_custom_fact_with_optional_target_item() { /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, &mut events, ) @@ -2409,6 +2419,7 @@ async fn item_review_summaries_do_not_cross_threads_with_reused_item_ids() { response: Box::new(sample_thread_start_response( "thread-2", /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, &mut events, ) @@ -2696,6 +2707,7 @@ async fn subagent_thread_started_inherits_parent_connection_for_new_thread() { /*ephemeral*/ false, "gpt-5", )), + thread_initialization: None, }, &mut events, ) @@ -3450,6 +3462,7 @@ async fn accepted_turn_steer_emits_expected_event() { connection_id: 7, request_id: RequestId::Integer(4), response: Box::new(sample_turn_steer_response("turn-2")), + thread_initialization: None, }, &mut out, ) @@ -3621,6 +3634,7 @@ async fn turn_start_error_response_discards_pending_start_request() { connection_id: 7, request_id: RequestId::Integer(3), response: Box::new(sample_turn_start_response("turn-2")), + thread_initialization: None, }, &mut out, ) @@ -3920,6 +3934,7 @@ async fn accepted_steers_increment_turn_steer_count() { connection_id: 7, request_id: RequestId::Integer(4), response: Box::new(sample_turn_steer_response("turn-2")), + thread_initialization: None, }, &mut out, ) @@ -3967,6 +3982,7 @@ async fn accepted_steers_increment_turn_steer_count() { connection_id: 7, request_id: RequestId::Integer(6), response: Box::new(sample_turn_steer_response("turn-2")), + thread_initialization: None, }, &mut out, ) diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index 6b8e1e7589..31500611bf 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -18,7 +18,6 @@ use crate::facts::SkillInvocation; use crate::facts::SkillInvokedInput; use crate::facts::SubAgentThreadStartedInput; use crate::facts::ThreadInitializationFact; -use crate::facts::ThreadInitializationInput; use crate::facts::TrackEventsContext; use crate::facts::TurnCodexErrorFact; use crate::facts::TurnProfileFact; @@ -319,6 +318,7 @@ impl AnalyticsEventsClient { connection_id: u64, request_id: RequestId, response: ClientResponsePayload, + thread_initialization: Option, ) { if !matches!( response, @@ -334,24 +334,10 @@ impl AnalyticsEventsClient { connection_id, request_id, response: Box::new(response), + thread_initialization, }); } - pub fn track_thread_initialization( - &self, - connection_id: u64, - request_id: RequestId, - fact: ThreadInitializationFact, - ) { - self.record_fact(AnalyticsFact::Custom( - CustomAnalyticsFact::ThreadInitialization(Box::new(ThreadInitializationInput { - connection_id, - request_id, - fact, - })), - )); - } - pub fn track_error_response( &self, connection_id: u64, diff --git a/codex-rs/analytics/src/client_tests.rs b/codex-rs/analytics/src/client_tests.rs index b7aa9f1c97..bfccecdcfc 100644 --- a/codex-rs/analytics/src/client_tests.rs +++ b/codex-rs/analytics/src/client_tests.rs @@ -252,7 +252,9 @@ fn track_response_only_enqueues_analytics_relevant_responses() { (RequestId::Integer(4), sample_turn_start_response()), (RequestId::Integer(5), sample_turn_steer_response()), ] { - client.track_response(/*connection_id*/ 7, request_id, response); + client.track_response( + /*connection_id*/ 7, request_id, response, /*thread_initialization*/ None, + ); assert!(matches!( receiver.try_recv(), Ok(AnalyticsFact::ClientResponse { .. }) @@ -263,6 +265,7 @@ fn track_response_only_enqueues_analytics_relevant_responses() { /*connection_id*/ 7, RequestId::Integer(6), ClientResponsePayload::ThreadArchive(ThreadArchiveResponse {}), + /*thread_initialization*/ None, ); assert!(matches!(receiver.try_recv(), Err(TryRecvError::Empty))); } diff --git a/codex-rs/analytics/src/facts.rs b/codex-rs/analytics/src/facts.rs index b5a9a43b99..64fda3f6ea 100644 --- a/codex-rs/analytics/src/facts.rs +++ b/codex-rs/analytics/src/facts.rs @@ -122,13 +122,6 @@ pub struct ThreadInitializationFact { pub profile: ThreadInitializationProfile, } -#[derive(Clone)] -pub(crate) struct ThreadInitializationInput { - pub connection_id: u64, - pub request_id: RequestId, - pub fact: ThreadInitializationFact, -} - #[derive(Clone)] pub struct TurnTokenUsageFact { pub turn_id: String, @@ -489,6 +482,7 @@ pub(crate) enum AnalyticsFact { connection_id: u64, request_id: RequestId, response: Box, + thread_initialization: Option, }, ErrorResponse { connection_id: u64, @@ -526,7 +520,6 @@ pub(crate) enum CustomAnalyticsFact { TurnResolvedConfig(Box), TurnTokenUsage(Box), TurnProfile(Box), - ThreadInitialization(Box), TurnCodexError(Box), SkillInvoked(SkillInvokedInput), AppMentioned(AppMentionedInput), diff --git a/codex-rs/analytics/src/reducer.rs b/codex-rs/analytics/src/reducer.rs index 704cfc3278..363c0b5886 100644 --- a/codex-rs/analytics/src/reducer.rs +++ b/codex-rs/analytics/src/reducer.rs @@ -70,7 +70,6 @@ use crate::facts::PluginUsedInput; use crate::facts::SkillInvokedInput; use crate::facts::SubAgentThreadStartedInput; use crate::facts::ThreadInitializationFact; -use crate::facts::ThreadInitializationInput; use crate::facts::ThreadInitializationMode; use crate::facts::ThreadInitializationProfile; use crate::facts::TurnCodexError; @@ -136,7 +135,6 @@ use std::path::PathBuf; #[derive(Default)] pub(crate) struct AnalyticsReducer { requests: HashMap<(u64, RequestId), RequestState>, - thread_initializations: HashMap<(u64, RequestId), ThreadInitializationFact>, turns: HashMap, connections: HashMap, threads: HashMap, @@ -410,9 +408,11 @@ impl AnalyticsReducer { connection_id, request_id, response, + thread_initialization, } => { if let Some(response) = response.into_client_response(request_id) { - self.ingest_response(connection_id, response, out).await; + self.ingest_response(connection_id, response, thread_initialization, out) + .await; } } AnalyticsFact::ErrorResponse { @@ -475,9 +475,6 @@ impl AnalyticsReducer { CustomAnalyticsFact::TurnProfile(input) => { self.ingest_turn_profile(*input, out).await; } - CustomAnalyticsFact::ThreadInitialization(input) => { - self.ingest_thread_initialization(*input); - } CustomAnalyticsFact::TurnCodexError(input) => { self.ingest_turn_codex_error(*input); } @@ -767,58 +764,39 @@ impl AnalyticsReducer { &mut self, connection_id: u64, response: ClientResponse, + thread_initialization: Option, out: &mut Vec, ) { match response { - ClientResponse::ThreadStart { - request_id, - response, - } => { - let initialization = self - .thread_initializations - .remove(&(connection_id, request_id)); + ClientResponse::ThreadStart { response, .. } => { self.emit_thread_initialized( connection_id, response.thread, response.model, - initialization + thread_initialization .map(|initialization| initialization.initialization_mode) .unwrap_or(ThreadInitializationMode::New), - initialization.map(|initialization| initialization.profile), + thread_initialization.map(|initialization| initialization.profile), out, ); } - ClientResponse::ThreadResume { - request_id, - response, - } => { - let profile = self - .thread_initializations - .remove(&(connection_id, request_id)) - .map(|initialization| initialization.profile); + ClientResponse::ThreadResume { response, .. } => { self.emit_thread_initialized( connection_id, response.thread, response.model, ThreadInitializationMode::Resumed, - profile, + thread_initialization.map(|initialization| initialization.profile), out, ); } - ClientResponse::ThreadFork { - request_id, - response, - } => { - let profile = self - .thread_initializations - .remove(&(connection_id, request_id)) - .map(|initialization| initialization.profile); + ClientResponse::ThreadFork { response, .. } => { self.emit_thread_initialized( connection_id, response.thread, response.model, ThreadInitializationMode::Forked, - profile, + thread_initialization.map(|initialization| initialization.profile), out, ); } @@ -848,11 +826,6 @@ impl AnalyticsReducer { } } - fn ingest_thread_initialization(&mut self, input: ThreadInitializationInput) { - self.thread_initializations - .insert((input.connection_id, input.request_id), input.fact); - } - fn ingest_server_request(&mut self, _connection_id: u64, request: ServerRequest) { match request { ServerRequest::CommandExecutionRequestApproval { request_id, params } => { diff --git a/codex-rs/app-server/src/outgoing_message.rs b/codex-rs/app-server/src/outgoing_message.rs index f28201b6cd..af20c8a2c3 100644 --- a/codex-rs/app-server/src/outgoing_message.rs +++ b/codex-rs/app-server/src/outgoing_message.rs @@ -526,20 +526,14 @@ impl OutgoingMessageSender { .into_jsonrpc_parts_and_payload(request_id.request_id.clone()) .map(|(id, result, response)| { if let Some(response) = response { - if let Some(thread_initialization) = request_context - .as_ref() - .and_then(|context| context.thread_initialization_timing.complete_request()) - { - self.analytics_events_client.track_thread_initialization( - connection_id.0, - request_id_for_analytics.clone(), - thread_initialization, - ); - } + let thread_initialization = request_context.as_ref().and_then(|context| { + context.thread_initialization_timing.complete_request() + }); self.analytics_events_client.track_response( connection_id.0, request_id_for_analytics, response, + thread_initialization, ); } (id, result)