From b3436d5a0a0f0d358d77ebce45f394a5ec47ef95 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 7 Jul 2026 22:03:08 +0300 Subject: [PATCH 1/3] test(app-server): cover external-agent cancellation and approval paths Add regression coverage for the live external-agent host and cancel RPC before /external-agent execution is ungated. These assert the current deny/reject gating so it cannot silently regress. Unit tests (thread_external_agent_processor): - terminal events (Completed/Failed/Cancelled) mark the run terminal so execute() does not double-emit Failed; progress events do not. - terminal events with no subscribers are dropped yet still recorded. - events fan out to every subscribed connection (late-event gating). - request_permission denies by default (RejectOnce). - perform_action rejects managed actions with the gating reason and emits a ProposedAction notification (action-routing-failure guard). - is_cancelled tracks the CancellationToken. Integration tests (suite/v2/thread_external_agent): - cancel rejects empty and unknown run ids with the expected messages. - an active run reaches a terminal event and is removed from the active-run registry so a follow-up cancel reports it inactive. --- .../thread_external_agent_processor.rs | 248 ++++++++++++++++++ .../tests/suite/v2/thread_external_agent.rs | 220 ++++++++++++++++ 2 files changed, 468 insertions(+) diff --git a/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs b/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs index f6ff0ed7c..52e7b336e 100644 --- a/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs @@ -1291,6 +1291,254 @@ mod tests { assert!(rx.try_recv().is_err()); } + #[tokio::test] + async fn external_agent_host_marks_terminal_events_and_leaves_progress_events_open() { + let (tx, _rx) = mpsc::channel(8); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000310") + .expect("thread id should parse"); + + for progress in [ + ThreadExternalAgentEvent::Status { + message: "working".to_string(), + }, + ThreadExternalAgentEvent::ProposedAction { + proposal: serde_json::json!({"type": "read-file"}), + }, + ] { + let host = AppServerExternalAgentHost::new( + outgoing.clone(), + ThreadStateManager::new(), + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + host.emit(progress).await; + assert!( + !host.terminal_sent(), + "progress events must not mark the run terminal" + ); + } + + for terminal in [ + ThreadExternalAgentEvent::Completed { + summary: Some("done".to_string()), + }, + ThreadExternalAgentEvent::Failed { + message: "boom".to_string(), + }, + ThreadExternalAgentEvent::Cancelled { reason: None }, + ] { + let host = AppServerExternalAgentHost::new( + outgoing.clone(), + ThreadStateManager::new(), + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + host.emit(terminal).await; + assert!( + host.terminal_sent(), + "terminal events must mark the run terminal even without subscribers" + ); + } + } + + #[tokio::test] + async fn external_agent_host_drops_terminal_events_without_subscribers() { + let (tx, mut rx) = mpsc::channel(4); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000311") + .expect("thread id should parse"); + let host = AppServerExternalAgentHost::new( + outgoing, + ThreadStateManager::new(), + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + + // A terminal event that arrives after every subscriber has gone away is + // dropped (never leaked to an unrelated connection) yet still records + // that the run finished so `execute` will not double-emit `Failed`. + host.emit(ThreadExternalAgentEvent::Cancelled { reason: None }) + .await; + + assert!(host.terminal_sent()); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn external_agent_host_delivers_events_to_every_subscribed_connection() { + let (tx, mut rx) = mpsc::channel(8); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_state_manager = ThreadStateManager::new(); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000312") + .expect("thread id should parse"); + for connection in [ConnectionId(7), ConnectionId(8)] { + thread_state_manager + .connection_initialized(connection, ConnectionCapabilities::default()) + .await; + assert!( + thread_state_manager + .try_add_connection_to_thread(thread_id, connection) + .await + ); + } + let host = AppServerExternalAgentHost::new( + outgoing, + thread_state_manager, + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + + host.emit(ThreadExternalAgentEvent::Status { + message: "working".to_string(), + }) + .await; + + let mut targeted = Vec::new(); + for _ in 0..2 { + let envelope = rx.recv().await.expect("targeted notification"); + let OutgoingEnvelope::ToConnection { connection_id, .. } = envelope else { + panic!("expected targeted envelope, got {envelope:?}"); + }; + targeted.push(connection_id); + } + targeted.sort_by_key(|connection| connection.0); + assert_eq!(targeted, vec![ConnectionId(7), ConnectionId(8)]); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn external_agent_host_denies_permission_requests_by_default() { + let (tx, _rx) = mpsc::channel(4); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000313") + .expect("thread id should parse"); + let host = AppServerExternalAgentHost::new( + outgoing, + ThreadStateManager::new(), + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + + let decision = host + .request_permission(ExternalAgentPermissionRequest { + id: "perm-1".to_string(), + action: ExternalAgentActionRequest::WriteFile { + path: PathBuf::from("/repo/file.txt"), + content: "data".to_string(), + }, + options: vec![ + ExternalAgentPermissionOption::AllowOnce, + ExternalAgentPermissionOption::RejectOnce, + ], + }) + .await + .expect("permission decision"); + + assert_eq!(decision, ExternalAgentPermissionDecision::RejectOnce); + } + + #[tokio::test] + async fn external_agent_host_rejects_managed_actions_and_notifies_proposal() { + let (tx, mut rx) = mpsc::channel(4); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_state_manager = ThreadStateManager::new(); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000314") + .expect("thread id should parse"); + thread_state_manager + .connection_initialized(ConnectionId(1), ConnectionCapabilities::default()) + .await; + assert!( + thread_state_manager + .try_add_connection_to_thread(thread_id, ConnectionId(1)) + .await + ); + let action = ExternalAgentActionRequest::RunCommand { + command: vec!["rm".to_string(), "-rf".to_string(), "/".to_string()], + cwd: None, + }; + let host = AppServerExternalAgentHost::new( + outgoing, + thread_state_manager, + thread_id, + "run-1".to_string(), + CancellationToken::new(), + ); + + let result = host + .perform_action(action.clone()) + .await + .expect("perform_action result"); + + assert_eq!( + result, + ExternalAgentActionResult::Rejected { + reason: "external-agent managed action routing is not enabled yet".to_string(), + } + ); + + let envelope = rx.recv().await.expect("proposed-action notification"); + let OutgoingEnvelope::ToConnection { message, .. } = envelope else { + panic!("expected targeted envelope, got {envelope:?}"); + }; + let OutgoingMessage::AppServerNotification(ServerNotification::ThreadExternalAgentEvent( + notification, + )) = message + else { + panic!("expected external-agent notification, got {message:?}"); + }; + assert_eq!( + notification.event, + ThreadExternalAgentEvent::ProposedAction { + proposal: action_json(&action), + } + ); + assert!(rx.try_recv().is_err()); + } + + #[tokio::test] + async fn external_agent_host_reports_cancellation_from_token() { + let (tx, _rx) = mpsc::channel(4); + let outgoing = Arc::new(OutgoingMessageSender::new( + tx, + codex_analytics::AnalyticsEventsClient::disabled(), + )); + let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000315") + .expect("thread id should parse"); + let token = CancellationToken::new(); + let host = AppServerExternalAgentHost::new( + outgoing, + ThreadStateManager::new(), + thread_id, + "run-1".to_string(), + token.clone(), + ); + + assert!(!host.is_cancelled().await); + token.cancel(); + assert!(host.is_cancelled().await); + } + #[test] fn external_agent_permission_profile_scopes_read_roots_with_network() { let cwd = tempfile::TempDir::new().expect("cwd tempdir"); diff --git a/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs b/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs index b9a0f55c6..65e1b0658 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs @@ -25,6 +25,7 @@ use std::path::Path; use std::time::Duration; use tempfile::TempDir; use tokio::time::timeout; +use wiremock::MockServer; const DEFAULT_TIMEOUT: Duration = Duration::from_secs(10); @@ -365,6 +366,225 @@ async fn thread_external_agent_permission_respond_unknown_request_is_not_accepte Ok(()) } +/// Boot an app-server whose PATH exposes a fake, long-sleeping `cursor-agent` +/// so live external-agent runs can be started and cancelled deterministically. +async fn start_ready_external_agent_server() -> Result<(McpProcess, MockServer, TempDir)> { + let tmp = TempDir::new()?; + let codex_home = tmp.path().join("codex_home"); + std::fs::create_dir(&codex_home)?; + + let server = create_mock_responses_server_sequence(vec![]).await; + write_mock_responses_config_toml( + codex_home.as_path(), + &server.uri(), + &BTreeMap::new(), + /*auto_compact_limit*/ 200_000, + /*requires_openai_auth*/ None, + "mock_provider", + "compact", + )?; + codex_login::save_auth_profile_metadata( + codex_home.as_path(), + "cursor-work", + codex_login::AuthProfileMetadata { + subscription_provider: codex_login::AuthProfileSubscriptionProvider::Cursor, + last_permissions: None, + }, + )?; + write_mock_provider_models_cache(codex_home.as_path())?; + let bin_dir = tmp.path().join("bin"); + std::fs::create_dir(&bin_dir)?; + write_fake_executable(&bin_dir, "agent")?; + write_fake_executable(&bin_dir, "cursor-agent")?; + let path = path_with_fake_bin(&bin_dir)?; + + let mut mcp = McpProcess::new_with_env( + codex_home.as_path(), + &[ + ("CODEWITH_AUTH_PROFILE", Some("cursor-work")), + ("PATH", Some(path.as_str())), + ], + ) + .await?; + timeout(DEFAULT_TIMEOUT, mcp.initialize()).await??; + Ok((mcp, server, tmp)) +} + +async fn start_thread(mcp: &mut McpProcess) -> Result { + let start_id = mcp + .send_thread_start_request(ThreadStartParams::default()) + .await?; + let start_resp: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(start_id)), + ) + .await??; + let ThreadStartResponse { thread, .. } = to_response(start_resp)?; + Ok(thread.id) +} + +#[tokio::test] +async fn thread_external_agent_cancel_rejects_empty_and_unknown_runs() -> Result<()> { + let (mut mcp, _server, _tmp) = start_ready_external_agent_server().await?; + let thread_id = start_thread(&mut mcp).await?; + + let empty_id = mcp + .send_thread_external_agent_cancel_request(ThreadExternalAgentCancelParams { + thread_id: thread_id.clone(), + run_id: " ".to_string(), + }) + .await?; + let empty_resp: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(empty_id)), + ) + .await??; + assert_eq!( + to_response::(empty_resp)?, + ThreadExternalAgentCancelResponse { + cancelled: false, + message: "runId must not be empty".to_string(), + } + ); + + let unknown_id = mcp + .send_thread_external_agent_cancel_request(ThreadExternalAgentCancelParams { + thread_id, + run_id: "ext_does_not_exist".to_string(), + }) + .await?; + let unknown_resp: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(unknown_id)), + ) + .await??; + assert_eq!( + to_response::(unknown_resp)?, + ThreadExternalAgentCancelResponse { + cancelled: false, + message: "external-agent run `ext_does_not_exist` is not active".to_string(), + } + ); + + Ok(()) +} + +// Live runs spawn a sandboxed child process, which is not available on Windows +// CI, so the cancellation-lifecycle assertions are gated to unix. +// +// This exercises the end-to-end lifecycle of a live run: RunStarted is emitted, +// cancellation is requested, the run reaches a terminal state and is removed +// from the active-run registry so a follow-up cancel reports it inactive. The +// terminal event is accepted as either `Cancelled` (when the sandboxed child +// hangs until the cancellation token fires) or `Failed` (when the platform +// sandbox cannot launch the child in this environment) so the registry-cleanup +// guarantee is asserted deterministically regardless of sandbox availability. +#[cfg(unix)] +#[tokio::test] +async fn thread_external_agent_cancel_clears_active_run_after_terminal_event() -> Result<()> { + let (mut mcp, _server, _tmp) = start_ready_external_agent_server().await?; + let thread_id = start_thread(&mut mcp).await?; + + let start_id = mcp + .send_thread_external_agent_start_request(ThreadExternalAgentStartParams { + thread_id: thread_id.clone(), + runtime_id: "cursor".to_string(), + task: "inspect the auth wiring".to_string(), + mode: ThreadExternalAgentMode::Plan, + }) + .await?; + let start_resp: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(start_id)), + ) + .await??; + let response: ThreadExternalAgentStartResponse = to_response(start_resp)?; + assert_eq!(response.status, ThreadExternalAgentStartStatus::Started); + let run_id = response.run_id.expect("external-agent run id"); + + // RunStarted is emitted before the fake agent is driven. + let started_notification = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_notification_message("thread/externalAgent/event"), + ) + .await??; + let started: ThreadExternalAgentEventNotification = + serde_json::from_value(started_notification.params.expect("event params"))?; + assert_eq!(started.run_id, run_id); + assert_eq!( + started.event, + ThreadExternalAgentEvent::RunStarted { + runtime_id: "cursor".to_string(), + mode: ThreadExternalAgentMode::Plan, + task: "inspect the auth wiring".to_string(), + } + ); + + // Request cancellation. Whether this wins the race against a fast sandbox + // failure is environment-dependent, so the acknowledgement boolean is not + // asserted here; the registry-cleanup poll below is the deterministic guard. + let cancel_id = mcp + .send_thread_external_agent_cancel_request(ThreadExternalAgentCancelParams { + thread_id: thread_id.clone(), + run_id: run_id.clone(), + }) + .await?; + let _cancel_resp: JSONRPCResponse = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(cancel_id)), + ) + .await??; + + // The run tears the child process down and emits a terminal event. + let terminal_notification = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_notification_message("thread/externalAgent/event"), + ) + .await??; + let terminal: ThreadExternalAgentEventNotification = + serde_json::from_value(terminal_notification.params.expect("event params"))?; + assert_eq!(terminal.run_id, run_id); + assert!( + matches!( + terminal.event, + ThreadExternalAgentEvent::Cancelled { .. } | ThreadExternalAgentEvent::Failed { .. } + ), + "expected a terminal Cancelled or Failed event, got: {:?}", + terminal.event + ); + + // Once the run finishes it is removed from the active registry, so a later + // cancel reports the run is no longer active. + let inactive_response: ThreadExternalAgentCancelResponse = timeout(DEFAULT_TIMEOUT, async { + loop { + let cancel_id = mcp + .send_thread_external_agent_cancel_request(ThreadExternalAgentCancelParams { + thread_id: thread_id.clone(), + run_id: run_id.clone(), + }) + .await?; + let cancel_resp: JSONRPCResponse = mcp + .read_stream_until_response_message(RequestId::Integer(cancel_id)) + .await?; + let response = to_response::(cancel_resp)?; + if !response.cancelled { + return Ok::(response); + } + tokio::time::sleep(Duration::from_millis(25)).await; + } + }) + .await??; + assert_eq!( + inactive_response, + ThreadExternalAgentCancelResponse { + cancelled: false, + message: format!("external-agent run `{run_id}` is not active"), + } + ); + + Ok(()) +} + fn path_with_fake_bin(bin_dir: &Path) -> Result { let existing_path = std::env::var("PATH").unwrap_or_default(); let path = std::env::join_paths( From e55eb119c4cbe3c751cf077154fa425eca7c32dd Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Thu, 9 Jul 2026 16:00:14 +0300 Subject: [PATCH 2/3] test(app-server): split external-agent test target --- codex-rs/app-server/tests/suite/v2/mod.rs | 1 - .../app-server/tests/suite/v2/thread_external_agent.rs | 10 +++++++--- codex-rs/app-server/tests/thread_external_agent.rs | 2 ++ 3 files changed, 9 insertions(+), 4 deletions(-) create mode 100644 codex-rs/app-server/tests/thread_external_agent.rs diff --git a/codex-rs/app-server/tests/suite/v2/mod.rs b/codex-rs/app-server/tests/suite/v2/mod.rs index 23daf0520..e33d43417 100644 --- a/codex-rs/app-server/tests/suite/v2/mod.rs +++ b/codex-rs/app-server/tests/suite/v2/mod.rs @@ -54,7 +54,6 @@ mod review; mod safety_check_downgrade; mod skills_list; mod thread_archive; -mod thread_external_agent; mod thread_fork; mod thread_inject_items; mod thread_list; diff --git a/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs b/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs index 65e1b0658..aab5b90cb 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_external_agent.rs @@ -587,9 +587,13 @@ async fn thread_external_agent_cancel_clears_active_run_after_terminal_event() - fn path_with_fake_bin(bin_dir: &Path) -> Result { let existing_path = std::env::var("PATH").unwrap_or_default(); - let path = std::env::join_paths( - std::iter::once(bin_dir.to_path_buf()).chain(std::env::split_paths(&existing_path)), - )?; + let path = std::env::join_paths(std::iter::once(bin_dir.to_path_buf()).chain( + std::env::split_paths(&existing_path).filter(|dir| { + !["claude", "claude.exe", "claude.cmd", "claude.bat"] + .iter() + .any(|program| dir.join(program).is_file()) + }), + ))?; Ok(path.to_string_lossy().into_owned()) } diff --git a/codex-rs/app-server/tests/thread_external_agent.rs b/codex-rs/app-server/tests/thread_external_agent.rs new file mode 100644 index 000000000..44813a775 --- /dev/null +++ b/codex-rs/app-server/tests/thread_external_agent.rs @@ -0,0 +1,2 @@ +#[path = "suite/v2/thread_external_agent.rs"] +mod thread_external_agent; From 34f985b4a2eaf2b9920ce0817a32e389c88a0926 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Thu, 9 Jul 2026 16:32:56 +0300 Subject: [PATCH 3/3] test(app-server): avoid external-agent permission test stall --- .../thread_external_agent_processor.rs | 76 +++++++++++++------ 1 file changed, 51 insertions(+), 25 deletions(-) diff --git a/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs b/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs index 52e7b336e..ac0f093fb 100644 --- a/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_external_agent_processor.rs @@ -1422,37 +1422,63 @@ mod tests { #[tokio::test] async fn external_agent_host_denies_permission_requests_by_default() { - let (tx, _rx) = mpsc::channel(4); - let outgoing = Arc::new(OutgoingMessageSender::new( - tx, - codex_analytics::AnalyticsEventsClient::disabled(), - )); let thread_id = ThreadId::from_string("00000000-0000-4000-8000-000000000313") .expect("thread id should parse"); - let host = AppServerExternalAgentHost::new( - outgoing, - ThreadStateManager::new(), - thread_id, - "run-1".to_string(), - CancellationToken::new(), - ); + let (host, mut rx) = subscribed_host(thread_id, "run-1", Duration::from_secs(30)).await; + let key = external_agent_permission_key("run-1", "perm-1"); + let thread_state = host.thread_state_manager.thread_state(thread_id).await; + + let await_host = host.clone(); + let handle = tokio::spawn(async move { + await_host + .request_permission(ExternalAgentPermissionRequest { + id: "perm-1".to_string(), + action: ExternalAgentActionRequest::WriteFile { + path: PathBuf::from("/repo/file.txt"), + content: "data".to_string(), + }, + options: vec![ + ExternalAgentPermissionOption::AllowOnce, + ExternalAgentPermissionOption::RejectOnce, + ], + }) + .await + }); - let decision = host - .request_permission(ExternalAgentPermissionRequest { - id: "perm-1".to_string(), - action: ExternalAgentActionRequest::WriteFile { - path: PathBuf::from("/repo/file.txt"), - content: "data".to_string(), + assert_eq!( + recv_external_agent_event(&mut rx).await, + ThreadExternalAgentEvent::PermissionRequested { + request: ThreadExternalAgentPermissionRequest { + id: "perm-1".to_string(), + action: action_json(&ExternalAgentActionRequest::WriteFile { + path: PathBuf::from("/repo/file.txt"), + content: "data".to_string(), + }), + options: vec![ + ThreadExternalAgentPermissionOption::AllowOnce, + ThreadExternalAgentPermissionOption::RejectOnce, + ], }, - options: vec![ - ExternalAgentPermissionOption::AllowOnce, - ExternalAgentPermissionOption::RejectOnce, - ], - }) - .await - .expect("permission decision"); + } + ); + assert!( + thread_state + .lock() + .await + .take_external_agent_permission(&key) + .is_some() + ); + let decision = handle.await.expect("join").expect("permission decision"); assert_eq!(decision, ExternalAgentPermissionDecision::RejectOnce); + assert_eq!( + recv_external_agent_event(&mut rx).await, + ThreadExternalAgentEvent::PermissionResolved { + request_id: "perm-1".to_string(), + decision: ThreadExternalAgentPermissionOption::RejectOnce, + resolution: ThreadExternalAgentPermissionResolution::DefaultDenied, + } + ); } #[tokio::test]