From ca3f48b04daaabfc6d176e8805e427d53138c269 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Tue, 24 Mar 2026 16:20:03 -0700 Subject: [PATCH] steer event --- codex-rs/core/src/analytics_client.rs | 87 +++++++++++ codex-rs/core/src/analytics_client_tests.rs | 31 ++++ codex-rs/core/src/codex.rs | 18 +++ codex-rs/core/tests/suite/items.rs | 156 ++++++++++++++++++++ 4 files changed, 292 insertions(+) diff --git a/codex-rs/core/src/analytics_client.rs b/codex-rs/core/src/analytics_client.rs index e5b5f13421..ed8465e7d9 100644 --- a/codex-rs/core/src/analytics_client.rs +++ b/codex-rs/core/src/analytics_client.rs @@ -46,6 +46,9 @@ pub(crate) struct CodexTurnEvent { pub(crate) num_input_images: usize, } +#[derive(Clone, Copy)] +pub(crate) struct CodexTurnSteerEvent; + pub(crate) fn build_track_events_context( model_slug: String, thread_id: String, @@ -110,6 +113,9 @@ impl AnalyticsEventsQueue { TrackEventsJob::TurnEvent(job) => { send_track_turn_event(&auth_manager, job).await; } + TrackEventsJob::TurnSteer(job) => { + send_track_turn_steer(&auth_manager, job).await; + } TrackEventsJob::PluginUsed(job) => { send_track_plugin_used(&auth_manager, job).await; } @@ -236,6 +242,19 @@ impl AnalyticsEventsClient { ); } + pub(crate) fn track_turn_steer( + &self, + tracking: TrackEventsContext, + turn_steer: CodexTurnSteerEvent, + ) { + track_turn_steer( + &self.queue, + Arc::clone(&self.config), + Some(tracking), + turn_steer, + ); + } + pub fn track_plugin_installed(&self, plugin: PluginTelemetryMetadata) { track_plugin_management( &self.queue, @@ -278,6 +297,7 @@ enum TrackEventsJob { AppMentioned(TrackAppMentionedJob), AppUsed(TrackAppUsedJob), TurnEvent(TrackTurnEventJob), + TurnSteer(TrackTurnSteerJob), PluginUsed(TrackPluginUsedJob), PluginInstalled(TrackPluginManagementJob), PluginUninstalled(TrackPluginManagementJob), @@ -309,6 +329,12 @@ struct TrackTurnEventJob { turn_event: CodexTurnEvent, } +struct TrackTurnSteerJob { + config: Arc, + tracking: TrackEventsContext, + turn_steer: CodexTurnSteerEvent, +} + struct TrackPluginUsedJob { config: Arc, tracking: TrackEventsContext, @@ -344,6 +370,7 @@ enum TrackEventRequest { AppMentioned(CodexAppMentionedEventRequest), AppUsed(CodexAppUsedEventRequest), TurnEvent(CodexTurnEventRequest), + TurnSteer(CodexTurnSteerEventRequest), PluginUsed(CodexPluginUsedEventRequest), PluginInstalled(CodexPluginEventRequest), PluginUninstalled(CodexPluginEventRequest), @@ -417,6 +444,20 @@ struct CodexTurnEventRequest { event_params: CodexTurnEventParams, } +#[derive(Serialize)] +struct CodexTurnSteerEventParams { + thread_id: String, + turn_id: String, + product_client_id: Option, + model: Option, +} + +#[derive(Serialize)] +struct CodexTurnSteerEventRequest { + event_type: &'static str, + event_params: CodexTurnSteerEventParams, +} + #[derive(Serialize)] struct CodexPluginMetadata { plugin_id: Option, @@ -538,6 +579,26 @@ pub(crate) fn track_turn_event( queue.try_send(job); } +pub(crate) fn track_turn_steer( + queue: &AnalyticsEventsQueue, + config: Arc, + tracking: Option, + turn_steer: CodexTurnSteerEvent, +) { + if config.analytics_enabled == Some(false) { + return; + } + let Some(tracking) = tracking else { + return; + }; + let job = TrackEventsJob::TurnSteer(TrackTurnSteerJob { + config, + tracking, + turn_steer, + }); + queue.try_send(job); +} + pub(crate) fn track_plugin_used( queue: &AnalyticsEventsQueue, config: Arc, @@ -677,6 +738,20 @@ async fn send_track_turn_event(auth_manager: &AuthManager, job: TrackTurnEventJo send_track_events(auth_manager, config, events).await; } +async fn send_track_turn_steer(auth_manager: &AuthManager, job: TrackTurnSteerJob) { + let TrackTurnSteerJob { + config, + tracking, + turn_steer, + } = job; + let events = vec![TrackEventRequest::TurnSteer(CodexTurnSteerEventRequest { + event_type: "codex_turn_event", + event_params: codex_turn_steer_event_params(&tracking, turn_steer), + })]; + + send_track_events(auth_manager, config, events).await; +} + async fn send_track_plugin_used(auth_manager: &AuthManager, job: TrackPluginUsedJob) { let TrackPluginUsedJob { config, @@ -767,6 +842,18 @@ fn codex_turn_event_params( } } +fn codex_turn_steer_event_params( + tracking: &TrackEventsContext, + _turn_steer: CodexTurnSteerEvent, +) -> CodexTurnSteerEventParams { + CodexTurnSteerEventParams { + thread_id: tracking.thread_id.clone(), + turn_id: tracking.turn_id.clone(), + product_client_id: Some(crate::default_client::originator().value), + model: Some(tracking.model_slug.clone()), + } +} + fn sandbox_policy_mode(sandbox_policy: &SandboxPolicy) -> &'static str { match sandbox_policy { SandboxPolicy::DangerFullAccess => "full_access", diff --git a/codex-rs/core/src/analytics_client_tests.rs b/codex-rs/core/src/analytics_client_tests.rs index 1a8cb80051..19e727a82d 100644 --- a/codex-rs/core/src/analytics_client_tests.rs +++ b/codex-rs/core/src/analytics_client_tests.rs @@ -6,6 +6,8 @@ use super::CodexPluginEventRequest; use super::CodexPluginUsedEventRequest; use super::CodexTurnEvent; use super::CodexTurnEventRequest; +use super::CodexTurnSteerEvent; +use super::CodexTurnSteerEventRequest; use super::InvocationType; use super::TrackEventRequest; use super::TrackEventsContext; @@ -13,6 +15,7 @@ use super::codex_app_metadata; use super::codex_plugin_metadata; use super::codex_plugin_used_metadata; use super::codex_turn_event_params; +use super::codex_turn_steer_event_params; use super::normalize_path_for_skill_id; use crate::plugins::AppConnectorId; use crate::plugins::PluginCapabilitySummary; @@ -250,6 +253,34 @@ fn turn_event_serializes_expected_shape() { ); } +#[test] +fn turn_steer_event_serializes_expected_shape() { + let tracking = TrackEventsContext { + model_slug: "gpt-5".to_string(), + thread_id: "thread-2".to_string(), + turn_id: "turn-2".to_string(), + }; + let event = TrackEventRequest::TurnSteer(CodexTurnSteerEventRequest { + event_type: "codex_turn_event", + event_params: codex_turn_steer_event_params(&tracking, CodexTurnSteerEvent), + }); + + let payload = serde_json::to_value(&event).expect("serialize turn steer event"); + + assert_eq!( + payload, + json!({ + "event_type": "codex_turn_event", + "event_params": { + "thread_id": "thread-2", + "turn_id": "turn-2", + "product_client_id": crate::default_client::originator().value, + "model": "gpt-5" + } + }) + ); +} + #[test] fn plugin_used_event_serializes_expected_shape() { let tracking = TrackEventsContext { diff --git a/codex-rs/core/src/codex.rs b/codex-rs/core/src/codex.rs index 5327771caa..4106ec07e2 100644 --- a/codex-rs/core/src/codex.rs +++ b/codex-rs/core/src/codex.rs @@ -15,6 +15,7 @@ use crate::agent::agent_status_from_event; use crate::analytics_client::AnalyticsEventsClient; use crate::analytics_client::AppInvocation; use crate::analytics_client::CodexTurnEvent; +use crate::analytics_client::CodexTurnSteerEvent; use crate::analytics_client::InvocationType; use crate::analytics_client::build_track_events_context; use crate::apps::render_apps_section; @@ -2739,6 +2740,16 @@ impl Session { )) } + async fn active_turn_tracking(&self) -> Option { + let active = self.active_turn.lock().await; + let (_, task) = active.as_ref()?.tasks.first()?; + Some(build_track_events_context( + task.turn_context.model_info.slug.clone(), + self.conversation_id.to_string(), + task.turn_context.sub_id.clone(), + )) + } + pub(crate) async fn record_execpolicy_amendment_message( &self, sub_id: &str, @@ -3913,6 +3924,13 @@ impl Session { let mut turn_state = active_turn.turn_state.lock().await; turn_state.push_pending_input(input.into()); + drop(turn_state); + drop(active); + if let Some(tracking) = self.active_turn_tracking().await { + self.services + .analytics_events_client + .track_turn_steer(tracking, CodexTurnSteerEvent); + } Ok(active_turn_id.clone()) } diff --git a/codex-rs/core/tests/suite/items.rs b/codex-rs/core/tests/suite/items.rs index f34fcca1de..f898727b32 100644 --- a/codex-rs/core/tests/suite/items.rs +++ b/codex-rs/core/tests/suite/items.rs @@ -33,6 +33,7 @@ use core_test_support::responses::ev_response_created; use core_test_support::responses::ev_web_search_call_added_partial; use core_test_support::responses::ev_web_search_call_done; use core_test_support::responses::mount_sse_once; +use core_test_support::responses::mount_sse_sequence; use core_test_support::responses::sse; use core_test_support::responses::start_mock_server; use core_test_support::skip_if_no_network; @@ -46,6 +47,7 @@ use std::path::Path; use std::path::PathBuf; use std::time::Duration; use std::time::Instant; +use tempfile::tempdir; fn image_generation_artifact_path(codex_home: &Path, session_id: &str, call_id: &str) -> PathBuf { fn sanitize(value: &str) -> String { @@ -71,6 +73,37 @@ fn image_generation_artifact_path(codex_home: &Path, session_id: &str, call_id: .join(format!("{}.png", sanitize(call_id))) } +async fn wait_for_analytics_event( + server: &wiremock::MockServer, + event_type: &str, + event_match: impl Fn(&Value) -> bool, +) -> Value { + let deadline = Instant::now() + Duration::from_secs(10); + loop { + let requests = server.received_requests().await.unwrap_or_default(); + if let Some(event) = requests + .into_iter() + .filter(|request| request.url.path() == "/codex/analytics-events/events") + .find_map(|request| { + let payload: Value = + serde_json::from_slice(&request.body).expect("analytics payload"); + payload["events"].as_array().and_then(|events| { + events + .iter() + .find(|event| event["event_type"] == event_type && event_match(event)) + .cloned() + }) + }) + { + return event; + } + if Instant::now() >= deadline { + panic!("timed out waiting for analytics event {event_type}"); + } + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn user_message_item_is_emitted() -> anyhow::Result<()> { skip_if_no_network!(Ok(())); @@ -237,6 +270,129 @@ async fn user_turn_tracks_turn_metadata_analytics() -> anyhow::Result<()> { Ok(()) } +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn user_turn_tracks_turn_steer_analytics() -> anyhow::Result<()> { + skip_if_no_network!(Ok(())); + + let server = start_mock_server().await; + let temp = tempdir()?; + let unblock_path = temp.path().join("unblock-steering"); + let command = format!( + "while [ ! -f \"{}\" ]; do sleep 0.01; done; echo done", + unblock_path.display() + ); + let call_id = "shell-steering-call"; + + mount_sse_sequence( + &server, + vec![ + sse(vec![ + ev_response_created("resp-1"), + core_test_support::responses::ev_function_call( + call_id, + "shell", + &serde_json::to_string(&serde_json::json!({ + "command": ["/bin/sh", "-c", command], + }))?, + ), + ev_completed("resp-1"), + ]), + sse(vec![ + ev_assistant_message("msg-2", "done"), + ev_completed("resp-2"), + ]), + ], + ) + .await; + + let chatgpt_base_url = server.uri(); + let test = test_codex() + .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) + .with_model("gpt-5") + .with_config(move |config| { + config.chatgpt_base_url = chatgpt_base_url; + }) + .build(&server) + .await?; + let codex = test.codex.clone(); + let turn_model = test.session_configured.model.clone(); + + codex + .submit(Op::UserTurn { + items: vec![UserInput::Text { + text: "start steering flow".into(), + text_elements: Vec::new(), + }], + final_output_json_schema: None, + cwd: test.cwd_path().to_path_buf(), + approval_policy: AskForApproval::Never, + approvals_reviewer: None, + sandbox_policy: SandboxPolicy::DangerFullAccess, + model: turn_model, + effort: None, + summary: None, + service_tier: None, + collaboration_mode: None, + personality: None, + }) + .await?; + + let turn_id = wait_for_event_match(&codex, |ev| match ev { + EventMsg::TurnStarted(event) => Some(event.turn_id.clone()), + _ => None, + }) + .await; + + wait_for_event_match(&codex, |ev| match ev { + EventMsg::ExecCommandBegin(event) if event.call_id == call_id => Some(()), + _ => None, + }) + .await; + + let steered_turn_id = codex + .steer_input( + vec![UserInput::Text { + text: "steering metadata check".into(), + text_elements: Vec::new(), + }], + Some(turn_id.as_str()), + ) + .await + .expect("steer should succeed on active turn"); + assert_eq!(steered_turn_id, turn_id); + + std::fs::write(&unblock_path, "go")?; + + wait_for_event(&codex, |event| matches!(event, EventMsg::TurnComplete(_))).await; + + let event = wait_for_analytics_event(&server, "codex_turn_event", |event| { + event["event_params"].get("model").is_some() + && event["event_params"].get("sandbox_policy").is_none() + }) + .await; + let event_params = &event["event_params"]; + + assert_eq!( + event_params["product_client_id"], + serde_json::json!(codex_core::default_client::originator().value) + ); + assert_eq!(event_params["model"], "gpt-5"); + assert!(event_params["thread_id"].as_str().is_some()); + assert!(event_params["turn_id"].as_str().is_some()); + assert!(event_params.get("sandbox_policy").is_none()); + assert!(event_params.get("reasoning_effort").is_none()); + assert!(event_params.get("reasoning_summary").is_none()); + assert!(event_params.get("service_tier").is_none()); + assert!(event_params.get("approval_policy").is_none()); + assert!(event_params.get("approvals_reviewer").is_none()); + assert!(event_params.get("sandbox_network_access").is_none()); + assert!(event_params.get("collaboration_mode").is_none()); + assert!(event_params.get("personality").is_none()); + assert!(event_params.get("num_input_images").is_none()); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn assistant_message_item_is_emitted() -> anyhow::Result<()> { skip_if_no_network!(Ok(()));