diff --git a/codex-rs/app-server/tests/suite/v2/analytics.rs b/codex-rs/app-server/tests/suite/v2/analytics.rs index 2acc86caee..c0ef692285 100644 --- a/codex-rs/app-server/tests/suite/v2/analytics.rs +++ b/codex-rs/app-server/tests/suite/v2/analytics.rs @@ -1,15 +1,37 @@ use anyhow::Result; use app_test_support::ChatGptAuthFixture; use app_test_support::DEFAULT_CLIENT_NAME; +use app_test_support::TestAppServer; +use app_test_support::create_final_assistant_message_sse_response; +use app_test_support::create_mock_responses_server_sequence; +use app_test_support::create_shell_command_sse_response; +use app_test_support::to_response; use app_test_support::write_chatgpt_auth; +use app_test_support::write_mock_responses_config_toml_with_chatgpt_base_url; +use codex_app_server_protocol::JSONRPCResponse; +use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::SandboxPolicy; +use codex_app_server_protocol::ThreadStartParams; +use codex_app_server_protocol::ThreadStartResponse; +use codex_app_server_protocol::TurnStartParams; +use codex_app_server_protocol::UserInput; use codex_config::types::AuthCredentialsStoreMode; use codex_config::types::OtelExporterKind; use codex_config::types::OtelHttpProtocol; use codex_core::config::ConfigBuilder; +use codex_core_plugins::loader::curated_plugin_cache_version; +use codex_core_plugins::store::PluginStore; +use codex_plugin::PluginId; +use core_test_support::responses; +use core_test_support::skip_if_no_network; +use core_test_support::skip_if_remote; +use core_test_support::skip_if_wine_exec; use pretty_assertions::assert_eq; use serde_json::Value; +use serde_json::json; use std::collections::HashMap; use std::path::Path; +use std::path::PathBuf; use std::time::Duration; use tempfile::TempDir; use tokio::time::timeout; @@ -230,3 +252,335 @@ pub(crate) fn assert_basic_thread_initialized_event( ); assert!(event["event_params"]["created_at"].as_u64().is_some()); } + +const METRICS_PLUGIN_ID: &str = "sample@openai-curated"; +const TEST_CURATED_PLUGIN_SHA: &str = "0123456789abcdef0123456789abcdef01234567"; + +#[derive(Clone, Copy)] +enum PluginMetricsRuntime { + Classic, + Unified, + UnifiedBackground, +} + +fn write_curated_metrics_plugin(codex_home: &Path) -> Result { + let plugin_id = PluginId::parse(METRICS_PLUGIN_ID)?; + let plugin_root = PluginStore::new(codex_home.to_path_buf()).plugin_root( + &plugin_id, + &curated_plugin_cache_version(TEST_CURATED_PLUGIN_SHA), + ); + let script_path = plugin_root.join("scripts/run.sh"); + std::fs::create_dir_all(plugin_root.join(".codex-plugin"))?; + std::fs::create_dir_all(script_path.parent().expect("script path has parent"))?; + std::fs::write( + plugin_root.join(".codex-plugin/plugin.json"), + r#"{"name":"sample","version":"0.1.0"}"#, + )?; + std::fs::write( + plugin_root.join("analytics.yaml"), + r#"version: 1 +operations: + scan: + path: ./scripts/run.sh + measurements: + findings: + dimensions: + severity: [high, low] + files_scanned: {} +"#, + )?; + std::fs::write( + &script_path, + r#"test -n "$CODEX_PLUGIN_METRICS_OUTPUT" +sleep "${1:-0.3}" +printf '%s' '{"version":1,"measurements":[{"name":"findings","value":3,"dimensions":{"severity":"high"}},{"name":"files_scanned","value":17}]}' > "$CODEX_PLUGIN_METRICS_OUTPUT" +"#, + )?; + + let curated_repo = codex_home.join(".tmp/plugins"); + std::fs::create_dir_all(curated_repo.join(".agents/plugins"))?; + std::fs::write( + curated_repo.join(".agents/plugins/marketplace.json"), + r#"{ + "name": "openai-curated", + "plugins": [ + {"name": "sample", "source": {"source": "local", "path": "./plugins/sample"}} + ] +}"#, + )?; + std::fs::write( + codex_home.join(".tmp/plugins.sha"), + format!("{TEST_CURATED_PLUGIN_SHA}\n"), + )?; + Ok(script_path.into_path_buf()) +} + +async fn assert_plugin_measurement_analytics(runtime: PluginMetricsRuntime) -> Result<()> { + skip_if_no_network!(Ok(())); + skip_if_remote!( + Ok(()), + "trusted plugin metrics fixture uses a local Codex home cache" + ); + skip_if_wine_exec!(Ok(()), "plugin metrics fixture is Unix-only"); + + let codex_home = TempDir::new()?; + let script_path = write_curated_metrics_plugin(codex_home.path())?; + let background = matches!(runtime, PluginMetricsRuntime::UnifiedBackground); + let unified_exec = !matches!(runtime, PluginMetricsRuntime::Classic); + let mut command = vec![ + "/bin/sh".to_string(), + script_path.to_string_lossy().into_owned(), + ]; + if background { + command.push("1.0".to_string()); + } + let call_id = "curated-plugin-metrics"; + let command_response = match runtime { + PluginMetricsRuntime::Classic => { + create_shell_command_sse_response(command, /*workdir*/ None, Some(5_000), call_id)? + } + PluginMetricsRuntime::Unified | PluginMetricsRuntime::UnifiedBackground => { + let arguments = serde_json::to_string(&json!({ + "cmd": shlex::try_join(command.iter().map(String::as_str))?, + "yield_time_ms": if background { 10 } else { 1_000 }, + }))?; + responses::sse(vec![ + responses::ev_response_created("resp-1"), + responses::ev_function_call(call_id, "exec_command", &arguments), + responses::ev_completed("resp-1"), + ]) + } + }; + let final_response = create_final_assistant_message_sse_response("done")?; + let server = + create_mock_responses_server_sequence(vec![command_response, final_response]).await; + + let analytics_server = responses::start_mock_server().await; + write_mock_responses_config_toml_with_chatgpt_base_url( + codex_home.path(), + &server.uri(), + &analytics_server.uri(), + )?; + let config_path = codex_home.path().join("config.toml"); + let config = std::fs::read_to_string(&config_path)?; + std::fs::write( + config_path, + format!( + r#"{config} +[features] +plugins = true +remote_plugin = false +unified_exec = {unified_exec} +shell_zsh_fork = false +unified_exec_zsh_fork = false + +[plugins."{METRICS_PLUGIN_ID}"] +enabled = true +"#, + ), + )?; + mount_analytics_capture(&analytics_server, codex_home.path()).await?; + + let mut mcp = TestAppServer::builder() + .with_codex_home(codex_home.path()) + .without_managed_config() + .build() + .await?; + timeout(Duration::from_secs(10), mcp.initialize()).await??; + let thread_request = mcp + .send_thread_start_request_with_auto_env(ThreadStartParams { + model: Some("mock-model".to_string()), + ..Default::default() + }) + .await?; + let thread_response: JSONRPCResponse = timeout( + Duration::from_secs(10), + mcp.read_stream_until_response_message(RequestId::Integer(thread_request)), + ) + .await??; + let ThreadStartResponse { thread, .. } = to_response(thread_response)?; + let thread_id = thread.id.clone(); + + let turn_request = mcp + .send_turn_start_request(TurnStartParams { + thread_id: thread.id, + input: vec![UserInput::Text { + text: "run the curated plugin metrics script".to_string(), + text_elements: Vec::new(), + }], + sandbox_policy: Some(SandboxPolicy::ReadOnly { + network_access: false, + }), + ..Default::default() + }) + .await?; + timeout( + Duration::from_secs(10), + mcp.read_stream_until_response_message(RequestId::Integer(turn_request)), + ) + .await??; + let completed_turn = timeout( + Duration::from_secs(10), + mcp.read_stream_until_notification_message("turn/completed"), + ) + .await??; + let turn_id = completed_turn + .params + .as_ref() + .and_then(|params| params["turn"]["id"].as_str()) + .expect("completed turn id"); + + if background { + let model_request_bodies = server + .received_requests() + .await + .unwrap_or_default() + .into_iter() + .filter(|request| request.url.path().ends_with("/responses")) + .map(|request| serde_json::from_slice::(&request.body)) + .collect::, _>>()?; + let request_text = serde_json::to_string(&model_request_bodies)?; + assert!(request_text.contains("Process running with session ID ")); + assert!(!request_text.contains("Process exited with code 0")); + } + + for measurement_name in ["findings", "files_scanned"] { + wait_for_matching_analytics_event(&analytics_server, Duration::from_secs(10), |event| { + event["event_type"] == "codex_plugin_measurement_event" + && event["event_params"]["item_id"] == call_id + && event["event_params"]["measurement_name"] == measurement_name + }) + .await?; + } + let command_event = + wait_for_matching_analytics_event(&analytics_server, Duration::from_secs(10), |event| { + event["event_type"] == "codex_command_execution_event" + && event["event_params"]["item_id"] == call_id + }) + .await?; + assert_eq!( + json!({ + "plugin_id": command_event["event_params"]["plugin_id"], + "script_path": command_event["event_params"]["script_path"], + "item_id": command_event["event_params"]["item_id"], + "exit_code": command_event["event_params"]["exit_code"], + }), + json!({ + "plugin_id": METRICS_PLUGIN_ID, + "script_path": "scripts/run.sh", + "item_id": call_id, + "exit_code": 0, + }) + ); + let mut measurement_events = Vec::new(); + for request in analytics_server + .received_requests() + .await + .unwrap_or_default() + { + if request.method != "POST" || request.url.path() != "/codex/analytics-events/events" { + continue; + } + let payload: Value = serde_json::from_slice(&request.body)?; + let Some(events) = payload["events"].as_array() else { + continue; + }; + measurement_events.extend( + events + .iter() + .filter(|event| { + event["event_type"] == "codex_plugin_measurement_event" + && event["event_params"]["item_id"] == call_id + }) + .cloned(), + ); + } + measurement_events.sort_by(|left, right| { + left["event_params"]["measurement_name"] + .as_str() + .cmp(&right["event_params"]["measurement_name"].as_str()) + }); + assert_eq!(measurement_events.len(), 2); + let execution_id = measurement_events[0]["event_params"]["execution_id"] + .as_str() + .expect("measurement execution id"); + assert!(!execution_id.is_empty()); + assert_eq!( + measurement_events[1]["event_params"]["execution_id"].as_str(), + Some(execution_id) + ); + assert_eq!( + measurement_events + .iter() + .map(|event| json!({ + "plugin_id": event["event_params"]["plugin_id"], + "operation": event["event_params"]["operation"], + "measurement_name": event["event_params"]["measurement_name"], + "number_value": event["event_params"]["number_value"], + "dimensions": event["event_params"]["dimensions"], + "item_id": event["event_params"]["item_id"], + })) + .collect::>(), + vec![ + json!({ + "plugin_id": METRICS_PLUGIN_ID, + "operation": "scan", + "measurement_name": "files_scanned", + "number_value": 17.0, + "dimensions": null, + "item_id": call_id, + }), + json!({ + "plugin_id": METRICS_PLUGIN_ID, + "operation": "scan", + "measurement_name": "findings", + "number_value": 3.0, + "dimensions": {"severity": "high"}, + "item_id": call_id, + }), + ] + ); + for event in &measurement_events { + let event_params = event["event_params"] + .as_object() + .expect("measurement event params"); + let mut field_names = event_params.keys().map(String::as_str).collect::>(); + field_names.sort_unstable(); + assert_eq!( + field_names, + vec![ + "dimensions", + "execution_id", + "item_id", + "measurement_name", + "number_value", + "operation", + "plugin_id", + "thread_id", + "turn_id", + ] + ); + assert_eq!(event_params["thread_id"], thread_id); + assert_eq!(event_params["turn_id"], turn_id); + } + + Ok(()) +} + +#[cfg_attr(windows, ignore = "plugin metrics fixture is Unix-only")] +#[tokio::test] +async fn classic_plugin_script_emits_measurement_analytics() -> Result<()> { + assert_plugin_measurement_analytics(PluginMetricsRuntime::Classic).await +} + +#[cfg_attr(windows, ignore = "plugin metrics fixture is Unix-only")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn unified_plugin_script_emits_measurement_analytics() -> Result<()> { + assert_plugin_measurement_analytics(PluginMetricsRuntime::Unified).await +} + +#[cfg_attr(windows, ignore = "plugin metrics fixture is Unix-only")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn unified_background_plugin_script_emits_measurements_after_turn_completion() -> Result<()> { + assert_plugin_measurement_analytics(PluginMetricsRuntime::UnifiedBackground).await +}