diff --git a/codex-rs/analytics/src/analytics_client_tests.rs b/codex-rs/analytics/src/analytics_client_tests.rs index bf97da4284..ab52079d7b 100644 --- a/codex-rs/analytics/src/analytics_client_tests.rs +++ b/codex-rs/analytics/src/analytics_client_tests.rs @@ -7,6 +7,8 @@ use crate::events::CodexCompactionEventRequest; use crate::events::CodexHookRunEventRequest; use crate::events::CodexPluginEventRequest; use crate::events::CodexPluginUsedEventRequest; +use crate::events::CodexResponsesApiCallEventParams; +use crate::events::CodexResponsesApiCallEventRequest; use crate::events::CodexRuntimeMetadata; use crate::events::CodexTurnEventRequest; use crate::events::GuardianApprovalRequestSource; @@ -29,6 +31,11 @@ use crate::facts::AppInvocation; use crate::facts::AppMentionedInput; use crate::facts::AppUsedInput; use crate::facts::CodexCompactionEvent; +use crate::facts::CodexResponseItemType; +use crate::facts::CodexResponsesApiCallFact; +use crate::facts::CodexResponsesApiCallStatus; +use crate::facts::CodexResponsesApiItemMetadata; +use crate::facts::CodexResponsesApiItemPhase; use crate::facts::CompactionImplementation; use crate::facts::CompactionPhase; use crate::facts::CompactionReason; @@ -91,6 +98,7 @@ use codex_plugin::PluginTelemetryMetadata; use codex_protocol::approvals::NetworkApprovalProtocol; use codex_protocol::config_types::ApprovalsReviewer; use codex_protocol::config_types::ModeKind; +use codex_protocol::models::MessagePhase; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::HookEventName; use codex_protocol::protocol::HookRunStatus; @@ -273,6 +281,25 @@ fn sample_turn_token_usage_fact(thread_id: &str, turn_id: &str) -> TurnTokenUsag } } +fn sample_responses_api_item_metadata( + item_phase: CodexResponsesApiItemPhase, + item_index: usize, +) -> CodexResponsesApiItemMetadata { + CodexResponsesApiItemMetadata { + item_phase, + item_index, + response_item_type: CodexResponseItemType::Message, + role: Some("assistant".to_string()), + status: Some("completed".to_string()), + message_phase: Some(MessagePhase::FinalAnswer), + call_id: None, + tool_name: None, + payload_bytes: Some(12), + text_part_count: Some(1), + image_part_count: Some(0), + } +} + fn sample_turn_completed_notification( thread_id: &str, turn_id: &str, @@ -752,6 +779,102 @@ fn compaction_event_serializes_expected_shape() { ); } +#[test] +fn responses_api_call_event_serializes_expected_shape() { + let event = TrackEventRequest::ResponsesApiCall(Box::new(CodexResponsesApiCallEventRequest { + event_type: "codex_responses_api_call_event", + event_params: CodexResponsesApiCallEventParams { + thread_id: "thread-1".to_string(), + turn_id: "turn-1".to_string(), + ephemeral: false, + thread_source: Some("user".to_string()), + initialization_mode: ThreadInitializationMode::New, + subagent_source: None, + parent_thread_id: None, + app_server_client: sample_app_server_client_metadata(), + runtime: sample_runtime_metadata(), + responses_id: Some("resp_123".to_string()), + turn_responses_call_index: 2, + model: Some("gpt-5".to_string()), + model_provider: Some("openai".to_string()), + reasoning_effort: Some("high".to_string()), + status: CodexResponsesApiCallStatus::Completed, + error: None, + started_at: 100, + completed_at: Some(102), + duration_ms: Some(2345), + input_item_count: 1, + output_item_count: 1, + input_tokens: Some(10), + cached_input_tokens: Some(4), + output_tokens: Some(20), + reasoning_output_tokens: Some(3), + total_tokens: Some(30), + items: vec![ + sample_responses_api_item_metadata( + CodexResponsesApiItemPhase::Input, + /*item_index*/ 0, + ), + sample_responses_api_item_metadata( + CodexResponsesApiItemPhase::Output, + /*item_index*/ 0, + ), + ], + }, + })); + + let payload = serde_json::to_value(&event).expect("serialize responses api call event"); + + assert_eq!(payload["event_type"], "codex_responses_api_call_event"); + assert_eq!(payload["event_params"]["thread_id"], "thread-1"); + assert_eq!(payload["event_params"]["turn_id"], "turn-1"); + assert_eq!(payload["event_params"]["ephemeral"], false); + assert_eq!(payload["event_params"]["thread_source"], "user"); + assert_eq!(payload["event_params"]["initialization_mode"], "new"); + assert_eq!( + payload["event_params"]["app_server_client"]["product_client_id"], + DEFAULT_ORIGINATOR + ); + assert_eq!( + payload["event_params"]["app_server_client"]["rpc_transport"], + "stdio" + ); + assert_eq!( + payload["event_params"]["runtime"]["codex_rs_version"], + "0.1.0" + ); + assert_eq!(payload["event_params"]["responses_id"], "resp_123"); + assert_eq!(payload["event_params"]["turn_responses_call_index"], 2); + assert_eq!(payload["event_params"]["model"], "gpt-5"); + assert_eq!(payload["event_params"]["model_provider"], "openai"); + assert_eq!(payload["event_params"]["reasoning_effort"], "high"); + assert_eq!(payload["event_params"]["status"], "completed"); + assert_eq!(payload["event_params"]["error"], json!(null)); + assert_eq!(payload["event_params"]["started_at"], 100); + assert_eq!(payload["event_params"]["completed_at"], 102); + assert_eq!(payload["event_params"]["duration_ms"], 2345); + assert_eq!(payload["event_params"]["input_item_count"], 1); + assert_eq!(payload["event_params"]["output_item_count"], 1); + assert_eq!(payload["event_params"]["input_tokens"], 10); + assert_eq!(payload["event_params"]["cached_input_tokens"], 4); + assert_eq!(payload["event_params"]["output_tokens"], 20); + assert_eq!(payload["event_params"]["reasoning_output_tokens"], 3); + assert_eq!(payload["event_params"]["total_tokens"], 30); + assert_eq!(payload["event_params"]["items"][0]["item_phase"], "input"); + assert_eq!( + payload["event_params"]["items"][0]["response_item_type"], + "message" + ); + assert_eq!(payload["event_params"]["items"][0]["role"], "assistant"); + assert_eq!( + payload["event_params"]["items"][0]["message_phase"], + "final_answer" + ); + assert_eq!(payload["event_params"]["items"][0]["text_part_count"], 1); + assert_eq!(payload["event_params"]["items"][0]["image_part_count"], 0); + assert_eq!(payload["event_params"]["items"][1]["item_phase"], "output"); +} + #[test] fn app_used_dedupe_is_keyed_by_turn_and_connector() { let (sender, _receiver) = mpsc::channel(1); @@ -1186,6 +1309,88 @@ async fn guardian_review_event_ingests_custom_fact_with_optional_target_item() { assert_eq!(payload[0]["event_params"]["review_timeout_ms"], 90_000); } +#[tokio::test] +async fn responses_api_call_event_ingests_custom_fact() { + let mut reducer = AnalyticsReducer::default(); + let mut events = Vec::new(); + + ingest_turn_prerequisites( + &mut reducer, + &mut events, + /*include_initialize*/ true, + /*include_resolved_config*/ true, + /*include_started*/ false, + /*include_token_usage*/ false, + ) + .await; + events.clear(); + + reducer + .ingest( + AnalyticsFact::Custom(CustomAnalyticsFact::ResponsesApiCall(Box::new( + CodexResponsesApiCallFact { + thread_id: "thread-2".to_string(), + turn_id: "turn-2".to_string(), + responses_id: Some("resp_456".to_string()), + turn_responses_call_index: 1, + status: CodexResponsesApiCallStatus::Completed, + error: None, + started_at: 100, + completed_at: Some(101), + duration_ms: Some(1200), + input_item_count: 1, + output_item_count: 1, + input_tokens: Some(123), + cached_input_tokens: Some(45), + output_tokens: Some(140), + reasoning_output_tokens: Some(13), + total_tokens: Some(321), + items: vec![ + sample_responses_api_item_metadata( + CodexResponsesApiItemPhase::Input, + /*item_index*/ 0, + ), + sample_responses_api_item_metadata( + CodexResponsesApiItemPhase::Output, + /*item_index*/ 0, + ), + ], + }, + ))), + &mut events, + ) + .await; + + let payload = serde_json::to_value(&events).expect("serialize events"); + assert_eq!(payload.as_array().expect("events array").len(), 1); + assert_eq!(payload[0]["event_type"], "codex_responses_api_call_event"); + assert_eq!(payload[0]["event_params"]["thread_id"], "thread-2"); + assert_eq!(payload[0]["event_params"]["turn_id"], "turn-2"); + assert_eq!(payload[0]["event_params"]["ephemeral"], false); + assert_eq!(payload[0]["event_params"]["thread_source"], "user"); + assert_eq!(payload[0]["event_params"]["initialization_mode"], "new"); + assert_eq!( + payload[0]["event_params"]["app_server_client"]["product_client_id"], + "codex-tui" + ); + assert_eq!(payload[0]["event_params"]["runtime"]["runtime_os"], "macos"); + assert_eq!(payload[0]["event_params"]["responses_id"], "resp_456"); + assert_eq!(payload[0]["event_params"]["turn_responses_call_index"], 1); + assert_eq!(payload[0]["event_params"]["model"], "gpt-5"); + assert_eq!(payload[0]["event_params"]["model_provider"], "openai"); + assert_eq!(payload[0]["event_params"]["reasoning_effort"], json!(null)); + assert_eq!(payload[0]["event_params"]["status"], "completed"); + assert_eq!(payload[0]["event_params"]["input_tokens"], 123); + assert_eq!( + payload[0]["event_params"]["items"][0]["item_phase"], + "input" + ); + assert_eq!( + payload[0]["event_params"]["items"][1]["item_phase"], + "output" + ); +} + #[test] fn subagent_thread_started_review_serializes_expected_shape() { let event = TrackEventRequest::ThreadInitialized(subagent_thread_started_event_request( diff --git a/codex-rs/analytics/src/client.rs b/codex-rs/analytics/src/client.rs index 25ab40953f..2f2a4a810d 100644 --- a/codex-rs/analytics/src/client.rs +++ b/codex-rs/analytics/src/client.rs @@ -8,6 +8,7 @@ use crate::facts::AnalyticsJsonRpcError; use crate::facts::AppInvocation; use crate::facts::AppMentionedInput; use crate::facts::AppUsedInput; +use crate::facts::CodexResponsesApiCallFact; use crate::facts::CustomAnalyticsFact; use crate::facts::HookRunFact; use crate::facts::HookRunInput; @@ -214,6 +215,12 @@ impl AnalyticsEventsClient { ))); } + pub fn track_responses_api_call(&self, event: CodexResponsesApiCallFact) { + self.record_fact(AnalyticsFact::Custom( + CustomAnalyticsFact::ResponsesApiCall(Box::new(event)), + )); + } + pub fn track_turn_resolved_config(&self, fact: TurnResolvedConfigFact) { self.record_fact(AnalyticsFact::Custom( CustomAnalyticsFact::TurnResolvedConfig(Box::new(fact)), diff --git a/codex-rs/analytics/src/events.rs b/codex-rs/analytics/src/events.rs index b542f8268f..8a1b7d2046 100644 --- a/codex-rs/analytics/src/events.rs +++ b/codex-rs/analytics/src/events.rs @@ -1,5 +1,7 @@ use crate::facts::AppInvocation; use crate::facts::CodexCompactionEvent; +use crate::facts::CodexResponsesApiCallStatus; +use crate::facts::CodexResponsesApiItemMetadata; use crate::facts::CompactionImplementation; use crate::facts::CompactionPhase; use crate::facts::CompactionReason; @@ -55,6 +57,8 @@ pub(crate) enum TrackEventRequest { AppUsed(CodexAppUsedEventRequest), HookRun(CodexHookRunEventRequest), Compaction(Box), + #[allow(dead_code)] + ResponsesApiCall(Box), TurnEvent(Box), TurnSteer(CodexTurnSteerEventRequest), PluginUsed(CodexPluginUsedEventRequest), @@ -311,6 +315,43 @@ pub(crate) struct CodexCompactionEventRequest { pub(crate) event_params: CodexCompactionEventParams, } +#[derive(Serialize)] +pub(crate) struct CodexResponsesApiCallEventParams { + pub(crate) thread_id: String, + pub(crate) turn_id: String, + pub(crate) ephemeral: bool, + pub(crate) thread_source: Option, + pub(crate) initialization_mode: ThreadInitializationMode, + pub(crate) subagent_source: Option, + pub(crate) parent_thread_id: Option, + pub(crate) app_server_client: CodexAppServerClientMetadata, + pub(crate) runtime: CodexRuntimeMetadata, + pub(crate) responses_id: Option, + pub(crate) turn_responses_call_index: u64, + pub(crate) model: Option, + pub(crate) model_provider: Option, + pub(crate) reasoning_effort: Option, + pub(crate) status: CodexResponsesApiCallStatus, + pub(crate) error: Option, + pub(crate) started_at: u64, + pub(crate) completed_at: Option, + pub(crate) duration_ms: Option, + pub(crate) input_item_count: usize, + pub(crate) output_item_count: usize, + pub(crate) input_tokens: Option, + pub(crate) cached_input_tokens: Option, + pub(crate) output_tokens: Option, + pub(crate) reasoning_output_tokens: Option, + pub(crate) total_tokens: Option, + pub(crate) items: Vec, +} + +#[derive(Serialize)] +pub(crate) struct CodexResponsesApiCallEventRequest { + pub(crate) event_type: &'static str, + pub(crate) event_params: CodexResponsesApiCallEventParams, +} + #[derive(Serialize)] pub(crate) struct CodexTurnEventParams { pub(crate) thread_id: String, diff --git a/codex-rs/analytics/src/facts.rs b/codex-rs/analytics/src/facts.rs index 1d371acb1c..10af798309 100644 --- a/codex-rs/analytics/src/facts.rs +++ b/codex-rs/analytics/src/facts.rs @@ -1,6 +1,10 @@ +use crate::byte_len; use crate::events::AppServerRpcTransport; use crate::events::CodexRuntimeMetadata; use crate::events::GuardianReviewEventParams; +use crate::nonzero_i64; +use crate::serialized_bytes; +use crate::serialized_string; use codex_app_server_protocol::ClientRequest; use codex_app_server_protocol::ClientResponse; use codex_app_server_protocol::InitializeParams; @@ -13,6 +17,14 @@ use codex_protocol::config_types::ModeKind; use codex_protocol::config_types::Personality; use codex_protocol::config_types::ReasoningSummary; use codex_protocol::config_types::ServiceTier; +use codex_protocol::models::ContentItem; +use codex_protocol::models::FunctionCallOutputBody; +use codex_protocol::models::FunctionCallOutputContentItem; +use codex_protocol::models::FunctionCallOutputPayload; +use codex_protocol::models::MessagePhase; +use codex_protocol::models::ReasoningItemContent; +use codex_protocol::models::ReasoningItemReasoningSummary; +use codex_protocol::models::ResponseItem; use codex_protocol::openai_models::ReasoningEffort; use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::HookEventName; @@ -262,6 +274,343 @@ pub struct CodexCompactionEvent { pub duration_ms: Option, } +#[derive(Clone, Copy, Debug, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum CodexResponsesApiCallStatus { + Completed, + Failed, + Interrupted, +} + +#[derive(Clone, Copy, Debug, Serialize, PartialEq, Eq)] +#[serde(rename_all = "lowercase")] +pub enum CodexResponsesApiItemPhase { + Input, + Output, +} + +#[derive(Clone, Copy, Debug, Serialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum CodexResponseItemType { + Message, + Reasoning, + LocalShellCall, + FunctionCall, + FunctionCallOutput, + CustomToolCall, + CustomToolCallOutput, + ToolSearchCall, + ToolSearchOutput, + WebSearchCall, + ImageGenerationCall, + GhostSnapshot, + Compaction, + Other, +} + +#[derive(Clone, Debug, Serialize, PartialEq, Eq)] +pub struct CodexResponsesApiItemMetadata { + pub item_phase: CodexResponsesApiItemPhase, + pub item_index: usize, + pub response_item_type: CodexResponseItemType, + pub role: Option, + pub status: Option, + pub message_phase: Option, + pub call_id: Option, + pub tool_name: Option, + pub payload_bytes: Option, + pub text_part_count: Option, + pub image_part_count: Option, +} + +#[derive(Clone, Debug)] +pub struct CodexResponsesApiCallFact { + pub thread_id: String, + pub turn_id: String, + pub responses_id: Option, + pub turn_responses_call_index: u64, + pub status: CodexResponsesApiCallStatus, + pub error: Option, + pub started_at: u64, + pub completed_at: Option, + pub duration_ms: Option, + pub input_item_count: usize, + pub output_item_count: usize, + pub input_tokens: Option, + pub cached_input_tokens: Option, + pub output_tokens: Option, + pub reasoning_output_tokens: Option, + pub total_tokens: Option, + pub items: Vec, +} + +pub fn response_items_metadata( + phase: CodexResponsesApiItemPhase, + items: &[ResponseItem], +) -> Vec { + items + .iter() + .enumerate() + .map(|(item_index, item)| response_item_metadata(phase, item_index, item)) + .collect() +} + +fn response_item_metadata( + item_phase: CodexResponsesApiItemPhase, + item_index: usize, + item: &ResponseItem, +) -> CodexResponsesApiItemMetadata { + let mut metadata = + CodexResponsesApiItemMetadata::new(item_phase, item_index, response_item_type(item)); + + match item { + ResponseItem::Message { + role, + content, + phase, + .. + } => { + metadata.role = Some(role.clone()); + metadata.message_phase = phase.clone(); + metadata.payload_bytes = nonzero_i64(message_content_text_bytes(content)); + let (text_part_count, image_part_count) = message_content_part_counts(content); + metadata.text_part_count = Some(text_part_count); + metadata.image_part_count = Some(image_part_count); + } + ResponseItem::Reasoning { + summary, + content, + encrypted_content, + .. + } => { + metadata.payload_bytes = encrypted_content + .as_ref() + .map(|value| byte_len(value)) + .or_else(|| nonzero_i64(reasoning_content_bytes(summary, content))); + metadata.text_part_count = + Some(summary.len() + content.as_ref().map(std::vec::Vec::len).unwrap_or_default()); + metadata.image_part_count = Some(0); + } + ResponseItem::LocalShellCall { + call_id, + status, + action, + .. + } => { + metadata.call_id = call_id.clone(); + metadata.tool_name = Some("local_shell".to_string()); + metadata.status = serialized_string(status); + metadata.payload_bytes = serialized_bytes(action); + } + ResponseItem::FunctionCall { + name, + arguments, + call_id, + .. + } => { + metadata.call_id = Some(call_id.clone()); + metadata.tool_name = Some(name.clone()); + metadata.payload_bytes = Some(byte_len(arguments)); + } + ResponseItem::ToolSearchCall { + call_id, + status, + arguments, + .. + } => { + metadata.call_id = call_id.clone(); + metadata.tool_name = Some("tool_search".to_string()); + metadata.status = status.clone(); + metadata.payload_bytes = serialized_bytes(arguments); + } + ResponseItem::FunctionCallOutput { call_id, output } => { + metadata.call_id = Some(call_id.clone()); + metadata.payload_bytes = function_call_output_bytes(output); + let (text_part_count, image_part_count) = function_call_output_part_counts(output); + metadata.text_part_count = text_part_count; + metadata.image_part_count = image_part_count; + } + ResponseItem::CustomToolCall { + status, + call_id, + name, + input, + .. + } => { + metadata.call_id = Some(call_id.clone()); + metadata.tool_name = Some(name.clone()); + metadata.status = status.clone(); + metadata.payload_bytes = Some(byte_len(input)); + } + ResponseItem::CustomToolCallOutput { + call_id, + name, + output, + } => { + metadata.call_id = Some(call_id.clone()); + metadata.tool_name = name.clone(); + metadata.payload_bytes = function_call_output_bytes(output); + let (text_part_count, image_part_count) = function_call_output_part_counts(output); + metadata.text_part_count = text_part_count; + metadata.image_part_count = image_part_count; + } + ResponseItem::ToolSearchOutput { + call_id, + status, + tools, + .. + } => { + metadata.call_id = call_id.clone(); + metadata.tool_name = Some("tool_search".to_string()); + metadata.status = Some(status.clone()); + metadata.payload_bytes = serialized_bytes(tools); + } + ResponseItem::WebSearchCall { status, action, .. } => { + metadata.tool_name = Some("web_search".to_string()); + metadata.status = status.clone(); + metadata.payload_bytes = action.as_ref().and_then(serialized_bytes); + } + ResponseItem::ImageGenerationCall { + id, + status, + revised_prompt, + result, + } => { + metadata.call_id = Some(id.clone()); + metadata.tool_name = Some("image_generation".to_string()); + metadata.status = Some(status.clone()); + metadata.payload_bytes = nonzero_i64(byte_len(result)) + .or_else(|| revised_prompt.as_ref().map(|value| byte_len(value))); + } + ResponseItem::Compaction { encrypted_content } => { + metadata.payload_bytes = Some(byte_len(encrypted_content)); + } + ResponseItem::GhostSnapshot { .. } | ResponseItem::Other => {} + } + + metadata +} + +impl CodexResponsesApiItemMetadata { + fn new( + item_phase: CodexResponsesApiItemPhase, + item_index: usize, + response_item_type: CodexResponseItemType, + ) -> Self { + Self { + item_phase, + item_index, + response_item_type, + role: None, + status: None, + message_phase: None, + call_id: None, + tool_name: None, + payload_bytes: None, + text_part_count: None, + image_part_count: None, + } + } +} + +fn response_item_type(item: &ResponseItem) -> CodexResponseItemType { + match item { + ResponseItem::Message { .. } => CodexResponseItemType::Message, + ResponseItem::Reasoning { .. } => CodexResponseItemType::Reasoning, + ResponseItem::LocalShellCall { .. } => CodexResponseItemType::LocalShellCall, + ResponseItem::FunctionCall { .. } => CodexResponseItemType::FunctionCall, + ResponseItem::ToolSearchCall { .. } => CodexResponseItemType::ToolSearchCall, + ResponseItem::FunctionCallOutput { .. } => CodexResponseItemType::FunctionCallOutput, + ResponseItem::CustomToolCall { .. } => CodexResponseItemType::CustomToolCall, + ResponseItem::CustomToolCallOutput { .. } => CodexResponseItemType::CustomToolCallOutput, + ResponseItem::ToolSearchOutput { .. } => CodexResponseItemType::ToolSearchOutput, + ResponseItem::WebSearchCall { .. } => CodexResponseItemType::WebSearchCall, + ResponseItem::ImageGenerationCall { .. } => CodexResponseItemType::ImageGenerationCall, + ResponseItem::GhostSnapshot { .. } => CodexResponseItemType::GhostSnapshot, + ResponseItem::Compaction { .. } => CodexResponseItemType::Compaction, + ResponseItem::Other => CodexResponseItemType::Other, + } +} + +fn message_content_text_bytes(content: &[ContentItem]) -> i64 { + content + .iter() + .map(|item| match item { + ContentItem::InputText { text } | ContentItem::OutputText { text } => byte_len(text), + ContentItem::InputImage { .. } => 0, + }) + .sum() +} + +fn message_content_part_counts(content: &[ContentItem]) -> (usize, usize) { + let mut text_part_count = 0; + let mut image_part_count = 0; + for item in content { + match item { + ContentItem::InputText { .. } | ContentItem::OutputText { .. } => { + text_part_count += 1; + } + ContentItem::InputImage { .. } => { + image_part_count += 1; + } + } + } + (text_part_count, image_part_count) +} + +fn reasoning_content_bytes( + summary: &[ReasoningItemReasoningSummary], + content: &Option>, +) -> i64 { + let summary_bytes = summary + .iter() + .map(|summary| match summary { + ReasoningItemReasoningSummary::SummaryText { text } => byte_len(text), + }) + .sum::(); + let content_bytes = content + .as_ref() + .map(|content| { + content + .iter() + .map(|content| match content { + ReasoningItemContent::ReasoningText { text } + | ReasoningItemContent::Text { text } => byte_len(text), + }) + .sum::() + }) + .unwrap_or_default(); + summary_bytes + content_bytes +} + +fn function_call_output_bytes(output: &FunctionCallOutputPayload) -> Option { + match &output.body { + FunctionCallOutputBody::Text(text) => Some(byte_len(text)), + FunctionCallOutputBody::ContentItems(items) => serialized_bytes(items), + } +} + +fn function_call_output_part_counts( + output: &FunctionCallOutputPayload, +) -> (Option, Option) { + let Some(content_items) = output.content_items() else { + return (None, None); + }; + let mut text_part_count = 0; + let mut image_part_count = 0; + for item in content_items { + match item { + FunctionCallOutputContentItem::InputText { .. } => { + text_part_count += 1; + } + FunctionCallOutputContentItem::InputImage { .. } => { + image_part_count += 1; + } + } + } + (Some(text_part_count), Some(image_part_count)) +} + #[allow(dead_code)] pub(crate) enum AnalyticsFact { Initialize { @@ -295,6 +644,7 @@ pub(crate) enum AnalyticsFact { pub(crate) enum CustomAnalyticsFact { SubAgentThreadStarted(SubAgentThreadStartedInput), Compaction(Box), + ResponsesApiCall(Box), GuardianReview(Box), TurnResolvedConfig(Box), TurnTokenUsage(Box), @@ -349,3 +699,70 @@ pub(crate) enum PluginState { Enabled, Disabled, } + +#[cfg(test)] +mod tests { + use super::*; + use codex_protocol::models::ImageDetail; + + #[test] + fn maps_message_metadata() { + let items = vec![ResponseItem::Message { + id: None, + role: "assistant".to_string(), + content: vec![ + ContentItem::OutputText { + text: "hello".to_string(), + }, + ContentItem::InputImage { + image_url: "data:image/png;base64,abc".to_string(), + detail: None, + }, + ], + end_turn: None, + phase: Some(MessagePhase::FinalAnswer), + }]; + + let metadata = response_items_metadata(CodexResponsesApiItemPhase::Output, &items); + + assert_eq!(metadata[0].item_phase, CodexResponsesApiItemPhase::Output); + assert_eq!( + metadata[0].response_item_type, + CodexResponseItemType::Message + ); + assert_eq!(metadata[0].role.as_deref(), Some("assistant")); + assert_eq!(metadata[0].message_phase, Some(MessagePhase::FinalAnswer)); + assert_eq!(metadata[0].payload_bytes, Some(5)); + assert_eq!(metadata[0].text_part_count, Some(1)); + assert_eq!(metadata[0].image_part_count, Some(1)); + } + + #[test] + fn maps_tool_call_output_metadata() { + let items = vec![ResponseItem::CustomToolCallOutput { + call_id: "call_1".to_string(), + name: Some("custom_tool".to_string()), + output: FunctionCallOutputPayload::from_content_items(vec![ + FunctionCallOutputContentItem::InputText { + text: "result".to_string(), + }, + FunctionCallOutputContentItem::InputImage { + image_url: "https://example.test/image.png".to_string(), + detail: Some(ImageDetail::High), + }, + ]), + }]; + + let metadata = response_items_metadata(CodexResponsesApiItemPhase::Output, &items); + + assert_eq!( + metadata[0].response_item_type, + CodexResponseItemType::CustomToolCallOutput + ); + assert_eq!(metadata[0].call_id.as_deref(), Some("call_1")); + assert_eq!(metadata[0].tool_name.as_deref(), Some("custom_tool")); + assert!(metadata[0].payload_bytes.unwrap_or_default() > 0); + assert_eq!(metadata[0].text_part_count, Some(1)); + assert_eq!(metadata[0].image_part_count, Some(1)); + } +} diff --git a/codex-rs/analytics/src/lib.rs b/codex-rs/analytics/src/lib.rs index 5c4cdfac79..bc0bac64b4 100644 --- a/codex-rs/analytics/src/lib.rs +++ b/codex-rs/analytics/src/lib.rs @@ -3,6 +3,7 @@ mod events; mod facts; mod reducer; +use serde::Serialize; use std::time::SystemTime; use std::time::UNIX_EPOCH; @@ -18,6 +19,11 @@ pub use events::GuardianReviewedAction; pub use facts::AnalyticsJsonRpcError; pub use facts::AppInvocation; pub use facts::CodexCompactionEvent; +pub use facts::CodexResponseItemType; +pub use facts::CodexResponsesApiCallFact; +pub use facts::CodexResponsesApiCallStatus; +pub use facts::CodexResponsesApiItemMetadata; +pub use facts::CodexResponsesApiItemPhase; pub use facts::CodexTurnSteerEvent; pub use facts::CompactionImplementation; pub use facts::CompactionPhase; @@ -39,6 +45,7 @@ pub use facts::TurnSteerRequestError; pub use facts::TurnSteerResult; pub use facts::TurnTokenUsageFact; pub use facts::build_track_events_context; +pub use facts::response_items_metadata; #[cfg(test)] mod analytics_client_tests; @@ -49,3 +56,24 @@ pub fn now_unix_seconds() -> u64 { .unwrap_or_default() .as_secs() } + +pub(crate) fn serialized_string(value: &T) -> Option { + match serde_json::to_value(value).ok()? { + serde_json::Value::String(value) => Some(value), + value => Some(value.to_string()), + } +} + +pub(crate) fn serialized_bytes(value: &T) -> Option { + serde_json::to_string(value) + .ok() + .map(|value| byte_len(&value)) +} + +pub(crate) fn nonzero_i64(value: i64) -> Option { + (value > 0).then_some(value) +} + +pub(crate) fn byte_len(value: &str) -> i64 { + i64::try_from(value.len()).unwrap_or(i64::MAX) +} diff --git a/codex-rs/analytics/src/reducer.rs b/codex-rs/analytics/src/reducer.rs index a6ce3fc831..4d70242d09 100644 --- a/codex-rs/analytics/src/reducer.rs +++ b/codex-rs/analytics/src/reducer.rs @@ -6,6 +6,8 @@ use crate::events::CodexCompactionEventRequest; use crate::events::CodexHookRunEventRequest; use crate::events::CodexPluginEventRequest; use crate::events::CodexPluginUsedEventRequest; +use crate::events::CodexResponsesApiCallEventParams; +use crate::events::CodexResponsesApiCallEventRequest; use crate::events::CodexRuntimeMetadata; use crate::events::CodexTurnEventParams; use crate::events::CodexTurnEventRequest; @@ -33,6 +35,7 @@ use crate::facts::AnalyticsJsonRpcError; use crate::facts::AppMentionedInput; use crate::facts::AppUsedInput; use crate::facts::CodexCompactionEvent; +use crate::facts::CodexResponsesApiCallFact; use crate::facts::CustomAnalyticsFact; use crate::facts::HookRunInput; use crate::facts::PluginState; @@ -202,6 +205,9 @@ impl AnalyticsReducer { CustomAnalyticsFact::Compaction(input) => { self.ingest_compaction(*input, out); } + CustomAnalyticsFact::ResponsesApiCall(input) => { + self.ingest_responses_api_call(*input, out); + } CustomAnalyticsFact::GuardianReview(input) => { self.ingest_guardian_review(*input, out); } @@ -739,6 +745,91 @@ impl AnalyticsReducer { ))); } + fn ingest_responses_api_call( + &mut self, + input: CodexResponsesApiCallFact, + out: &mut Vec, + ) { + let Some(connection_id) = self.thread_connections.get(&input.thread_id) else { + tracing::warn!( + thread_id = %input.thread_id, + turn_id = %input.turn_id, + "dropping responses api call analytics event: missing thread connection metadata" + ); + return; + }; + let Some(connection_state) = self.connections.get(connection_id) else { + tracing::warn!( + thread_id = %input.thread_id, + turn_id = %input.turn_id, + connection_id, + "dropping responses api call analytics event: missing connection metadata" + ); + return; + }; + let Some(thread_metadata) = self.thread_metadata.get(&input.thread_id) else { + tracing::warn!( + thread_id = %input.thread_id, + turn_id = %input.turn_id, + "dropping responses api call analytics event: missing thread lifecycle metadata" + ); + return; + }; + let Some(turn_state) = self.turns.get(&input.turn_id) else { + tracing::warn!( + thread_id = %input.thread_id, + turn_id = %input.turn_id, + "dropping responses api call analytics event: missing turn metadata" + ); + return; + }; + let Some(resolved_config) = turn_state.resolved_config.as_ref() else { + tracing::warn!( + thread_id = %input.thread_id, + turn_id = %input.turn_id, + "dropping responses api call analytics event: missing turn resolved config" + ); + return; + }; + + out.push(TrackEventRequest::ResponsesApiCall(Box::new( + CodexResponsesApiCallEventRequest { + event_type: "codex_responses_api_call_event", + event_params: CodexResponsesApiCallEventParams { + thread_id: input.thread_id, + turn_id: input.turn_id, + ephemeral: resolved_config.ephemeral, + thread_source: thread_metadata.thread_source.map(str::to_string), + initialization_mode: thread_metadata.initialization_mode, + subagent_source: thread_metadata.subagent_source.clone(), + parent_thread_id: thread_metadata.parent_thread_id.clone(), + app_server_client: connection_state.app_server_client.clone(), + runtime: connection_state.runtime.clone(), + responses_id: input.responses_id, + turn_responses_call_index: input.turn_responses_call_index, + model: Some(resolved_config.model.clone()), + model_provider: Some(resolved_config.model_provider.clone()), + reasoning_effort: resolved_config + .reasoning_effort + .map(|value| value.to_string()), + status: input.status, + error: input.error, + started_at: input.started_at, + completed_at: input.completed_at, + duration_ms: input.duration_ms, + input_item_count: input.input_item_count, + output_item_count: input.output_item_count, + input_tokens: input.input_tokens, + cached_input_tokens: input.cached_input_tokens, + output_tokens: input.output_tokens, + reasoning_output_tokens: input.reasoning_output_tokens, + total_tokens: input.total_tokens, + items: input.items, + }, + }, + ))); + } + fn ingest_turn_steer_response( &mut self, connection_id: u64, diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index efbcca6716..aad57dfab2 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -1,6 +1,9 @@ use std::collections::HashMap; use std::collections::HashSet; use std::sync::Arc; +use std::time::Instant; +use std::time::SystemTime; +use std::time::UNIX_EPOCH; use crate::SkillInjections; use crate::SkillLoadOutcome; @@ -58,11 +61,15 @@ use crate::unavailable_tool::collect_unavailable_called_tools; use crate::util::backoff; use crate::util::error_or_panic; use codex_analytics::AppInvocation; +use codex_analytics::CodexResponsesApiCallFact; +use codex_analytics::CodexResponsesApiCallStatus; +use codex_analytics::CodexResponsesApiItemPhase; use codex_analytics::CompactionPhase; use codex_analytics::CompactionReason; use codex_analytics::InvocationType; use codex_analytics::TurnResolvedConfigFact; use codex_analytics::build_track_events_context; +use codex_analytics::response_items_metadata; use codex_async_utils::OrCancelExt; use codex_features::Feature; use codex_hooks::HookEvent; @@ -90,6 +97,7 @@ use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::PlanDeltaEvent; use codex_protocol::protocol::ReasoningContentDeltaEvent; use codex_protocol::protocol::ReasoningRawContentDeltaEvent; +use codex_protocol::protocol::TokenUsage; use codex_protocol::protocol::TurnDiffEvent; use codex_protocol::protocol::WarningEvent; use codex_protocol::user_input::UserInput; @@ -394,6 +402,7 @@ pub(crate) async fn run_turn( // 1. At the start of a turn, so the fresh user prompt in `input` gets sampled first. // 2. After auto-compact, when model/tool continuation needs to resume before any steer. let mut can_drain_pending_input = input.is_empty(); + let mut next_turn_responses_call_index: u64 = 0; loop { if run_pending_session_start_hooks(&sess, &turn_context).await { @@ -465,12 +474,15 @@ pub(crate) async fn run_turn( .map(|user_message| user_message.message()) .collect::>(); let turn_metadata_header = turn_context.turn_metadata_state.current_header_value(); + let turn_responses_call_index = next_turn_responses_call_index; + next_turn_responses_call_index = next_turn_responses_call_index.saturating_add(1); match run_sampling_request( Arc::clone(&sess), Arc::clone(&turn_context), Arc::clone(&turn_diff_tracker), &mut client_session, turn_metadata_header.as_deref(), + turn_responses_call_index, sampling_request_input, &explicitly_enabled_connectors, skills_outcome, @@ -991,6 +1003,127 @@ pub(crate) fn build_prompt( } } +const RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES: usize = 1024; + +struct ResponsesApiCallAnalyticsState { + turn_responses_call_index: u64, + started_at: u64, + started_instant: Instant, + input_items: Vec, + output_items: Vec, + responses_id: Option, + token_usage: Option, +} + +impl ResponsesApiCallAnalyticsState { + fn new(turn_responses_call_index: u64) -> Self { + Self { + turn_responses_call_index, + started_at: current_unix_seconds(), + started_instant: Instant::now(), + input_items: Vec::new(), + output_items: Vec::new(), + responses_id: None, + token_usage: None, + } + } + + fn record_prompt(&mut self, prompt: &Prompt) { + self.input_items = prompt.input.clone(); + self.output_items.clear(); + self.responses_id = None; + self.token_usage = None; + } + + fn record_output_item_done(&mut self, item: &ResponseItem) { + self.output_items.push(item.clone()); + } + + fn record_completed(&mut self, response_id: String, token_usage: Option) { + self.responses_id = Some(response_id); + self.token_usage = token_usage; + } + + fn into_fact( + self, + sess: &Session, + turn_context: &TurnContext, + status: CodexResponsesApiCallStatus, + error: Option, + ) -> CodexResponsesApiCallFact { + let completed_at = current_unix_seconds(); + let mut items = + response_items_metadata(CodexResponsesApiItemPhase::Input, &self.input_items); + items.extend(response_items_metadata( + CodexResponsesApiItemPhase::Output, + &self.output_items, + )); + let token_usage = self.token_usage; + CodexResponsesApiCallFact { + thread_id: sess.conversation_id.to_string(), + turn_id: turn_context.sub_id.clone(), + responses_id: self.responses_id, + turn_responses_call_index: self.turn_responses_call_index, + status, + error, + started_at: self.started_at, + completed_at: Some(completed_at), + duration_ms: Some(self.started_instant.elapsed().as_millis() as u64), + input_item_count: self.input_items.len(), + output_item_count: self.output_items.len(), + input_tokens: token_usage.as_ref().map(|usage| usage.input_tokens), + cached_input_tokens: token_usage.as_ref().map(|usage| usage.cached_input_tokens), + output_tokens: token_usage.as_ref().map(|usage| usage.output_tokens), + reasoning_output_tokens: token_usage + .as_ref() + .map(|usage| usage.reasoning_output_tokens), + total_tokens: token_usage.as_ref().map(|usage| usage.total_tokens), + items, + } + } +} + +fn current_unix_seconds() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_secs() +} + +fn truncate_responses_api_error(mut error: String) -> String { + if error.len() <= RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES { + return error; + } + let mut truncate_at = RESPONSES_API_ANALYTICS_ERROR_MAX_BYTES; + while !error.is_char_boundary(truncate_at) { + truncate_at -= 1; + } + error.truncate(truncate_at); + error +} + +fn responses_api_call_error_status(err: &CodexErr) -> CodexResponsesApiCallStatus { + match err { + CodexErr::TurnAborted => CodexResponsesApiCallStatus::Interrupted, + _ => CodexResponsesApiCallStatus::Failed, + } +} + +fn track_responses_api_call_analytics( + sess: &Session, + turn_context: &TurnContext, + analytics: ResponsesApiCallAnalyticsState, + status: CodexResponsesApiCallStatus, + error: Option, +) { + if !sess.enabled(Feature::GeneralAnalytics) { + return; + } + sess.services + .analytics_events_client + .track_responses_api_call(analytics.into_fact(sess, turn_context, status, error)); +} + #[allow(clippy::too_many_arguments)] #[instrument(level = "trace", skip_all, @@ -1006,6 +1139,7 @@ async fn run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, client_session: &mut ModelClientSession, turn_metadata_header: Option<&str>, + turn_responses_call_index: u64, input: Vec, explicitly_enabled_connectors: &HashSet, skills_outcome: Option<&SkillLoadOutcome>, @@ -1042,6 +1176,8 @@ async fn run_sampling_request( .await; let mut retries = 0; let mut initial_input = Some(input); + let mut responses_api_call_analytics = + ResponsesApiCallAnalyticsState::new(turn_responses_call_index); loop { let prompt_input = if let Some(input) = initial_input.take() { input @@ -1056,6 +1192,7 @@ async fn run_sampling_request( turn_context.as_ref(), base_instructions.clone(), ); + responses_api_call_analytics.record_prompt(&prompt); let err = match try_run_sampling_request( tool_runtime.clone(), Arc::clone(&sess), @@ -1065,15 +1202,30 @@ async fn run_sampling_request( Arc::clone(&turn_diff_tracker), server_model_warning_emitted_for_turn, &prompt, + &mut responses_api_call_analytics, cancellation_token.child_token(), ) .await { Ok(output) => { + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Completed, + /*error*/ None, + ); return Ok(output); } Err(CodexErr::ContextWindowExceeded) => { sess.set_total_tokens_full(&turn_context).await; + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + Some("context window exceeded".to_string()), + ); return Err(CodexErr::ContextWindowExceeded); } Err(CodexErr::UsageLimitReached(e)) => { @@ -1081,12 +1233,28 @@ async fn run_sampling_request( if let Some(rate_limits) = rate_limits { sess.update_rate_limits(&turn_context, *rate_limits).await; } + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + Some(truncate_responses_api_error(e.to_string())), + ); return Err(CodexErr::UsageLimitReached(e)); } Err(err) => err, }; if !err.is_retryable() { + let status = responses_api_call_error_status(&err); + let error = Some(truncate_responses_api_error(format!("{err:#}"))); + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + status, + error, + ); return Err(err); } @@ -1138,6 +1306,14 @@ async fn run_sampling_request( } tokio::time::sleep(delay).await; } else { + let error = Some(truncate_responses_api_error(format!("{err:#}"))); + track_responses_api_call_analytics( + sess.as_ref(), + turn_context.as_ref(), + responses_api_call_analytics, + CodexResponsesApiCallStatus::Failed, + error, + ); return Err(err); } } @@ -1850,6 +2026,7 @@ async fn try_run_sampling_request( turn_diff_tracker: SharedTurnDiffTracker, server_model_warning_emitted_for_turn: &mut bool, prompt: &Prompt, + responses_api_call_analytics: &mut ResponsesApiCallAnalyticsState, cancellation_token: CancellationToken, ) -> CodexResult { feedback_tags!( @@ -1925,6 +2102,7 @@ async fn try_run_sampling_request( match event { ResponseEvent::Created => {} ResponseEvent::OutputItemDone(item) => { + responses_api_call_analytics.record_output_item_done(&item); if let Some((_, mut consumer)) = active_tool_argument_diff_consumer.take() && let Some(event) = consumer.flush_on_complete() { @@ -2095,9 +2273,10 @@ async fn try_run_sampling_request( sess.services.models_manager.refresh_if_new_etag(etag).await; } ResponseEvent::Completed { - response_id: _, + response_id, token_usage, } => { + responses_api_call_analytics.record_completed(response_id, token_usage.clone()); flush_assistant_text_segments_all( &sess, &turn_context,