Merge dispatcher lifecycle tests into client stack

This commit is contained in:
Chris Bookholt
2026-07-02 19:01:36 -07:00

View File

@@ -2224,6 +2224,7 @@ mod tests {
use codex_app_server_protocol::AutoReviewDecisionSource;
use codex_app_server_protocol::GuardianApprovalReviewStatus;
use codex_app_server_protocol::JSONRPCErrorError;
use codex_app_server_protocol::ServerRequest;
use codex_app_server_protocol::TurnPlanStepStatus;
use codex_login::CodexAuth;
use codex_protocol::AgentPath;
@@ -2386,6 +2387,100 @@ mod tests {
}
}
struct CommandApprovalTestContext {
_codex_home: TempDir,
conversation_id: ThreadId,
conversation: Arc<CodexThread>,
thread_manager: Arc<ThreadManager>,
outgoing: ThreadScopedOutgoingMessageSender,
raw_outgoing: Arc<OutgoingMessageSender>,
thread_state: Arc<Mutex<ThreadState>>,
thread_watch_manager: ThreadWatchManager,
rx: mpsc::Receiver<OutgoingEnvelope>,
}
impl CommandApprovalTestContext {
async fn new() -> Result<Self> {
let codex_home = TempDir::new()?;
let config = load_default_config_for_test(&codex_home).await;
let thread_manager = Arc::new(
codex_core::test_support::thread_manager_with_models_provider_and_home(
CodexAuth::create_dummy_chatgpt_auth_for_testing(),
config.model_provider.clone(),
config.codex_home.to_path_buf(),
Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()),
),
);
let codex_core::NewThread {
thread_id: conversation_id,
thread: conversation,
..
} = thread_manager.start_thread(config).await?;
let (tx, rx) = mpsc::channel(CHANNEL_CAPACITY);
let raw_outgoing = Arc::new(OutgoingMessageSender::new(
tx,
codex_analytics::AnalyticsEventsClient::disabled(),
));
let outgoing = ThreadScopedOutgoingMessageSender::new(
Arc::clone(&raw_outgoing),
vec![ConnectionId(1)],
conversation_id,
);
Ok(Self {
_codex_home: codex_home,
conversation_id,
conversation,
thread_manager,
outgoing,
raw_outgoing,
thread_state: new_thread_state(),
thread_watch_manager: ThreadWatchManager::new(),
rx,
})
}
async fn apply(&self, event: Event) {
apply_bespoke_event_handling(
event,
self.conversation_id,
Arc::clone(&self.conversation),
Arc::clone(&self.thread_manager),
self.outgoing.clone(),
Arc::clone(&self.thread_state),
self.thread_watch_manager.clone(),
Arc::new(tokio::sync::Semaphore::new(/*permits*/ 1)),
"test-provider".to_string(),
)
.await;
}
}
fn exec_approval_event(
call_id: &str,
approval_id: Option<&str>,
reason: Option<&str>,
) -> Event {
Event {
id: "turn-1".to_string(),
msg: EventMsg::ExecApprovalRequest(ExecApprovalRequestEvent {
call_id: call_id.to_string(),
approval_id: approval_id.map(str::to_string),
turn_id: "turn-1".to_string(),
environment_id: Some("local".to_string()),
started_at_ms: 1_000,
command: vec!["printf".to_string(), "hi".to_string()],
cwd: test_path_buf("/tmp").abs(),
reason: reason.map(str::to_string),
network_approval_context: None,
proposed_execpolicy_amendment: None,
proposed_network_policy_amendments: None,
additional_permissions: None,
available_decisions: Some(vec![ReviewDecision::Approved, ReviewDecision::Abort]),
parsed_cmd: Vec::new(),
}),
}
}
fn guardian_command_assessment(
id: &str,
turn_id: &str,
@@ -2610,266 +2705,188 @@ mod tests {
}
#[tokio::test]
async fn command_execution_started_helper_emits_once() -> Result<()> {
let conversation_id = ThreadId::new();
let thread_state = new_thread_state();
let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY);
let outgoing = Arc::new(OutgoingMessageSender::new(
tx,
codex_analytics::AnalyticsEventsClient::disabled(),
async fn repeated_approval_dispatch_preserves_item_and_callback_ids() -> Result<()> {
let mut context = CommandApprovalTestContext::new().await?;
context
.apply(exec_approval_event("cmd-1", None, Some("initial")))
.await;
let started = recv_broadcast_message(&mut context.rx).await?;
let OutgoingMessage::AppServerNotification(ServerNotification::ItemStarted(started)) =
started
else {
bail!("expected item/started");
};
assert!(matches!(
started.item,
ThreadItem::CommandExecution { ref id, .. } if id == "cmd-1"
));
let outgoing = ThreadScopedOutgoingMessageSender::new(
outgoing,
vec![ConnectionId(1)],
ThreadId::new(),
let initial_request = recv_broadcast_message(&mut context.rx).await?;
let OutgoingMessage::Request(ServerRequest::CommandExecutionRequestApproval {
request_id,
params,
}) = initial_request
else {
bail!("expected initial command approval request");
};
assert_eq!(params.item_id, "cmd-1");
assert_eq!(params.approval_id, None);
context
.raw_outgoing
.notify_client_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await;
context
.apply(exec_approval_event(
"cmd-1",
Some("retry-1"),
Some("sandbox denied"),
))
.await;
let retry_request = recv_broadcast_message(&mut context.rx).await?;
let OutgoingMessage::Request(ServerRequest::CommandExecutionRequestApproval {
request_id,
params,
}) = retry_request
else {
bail!("retry must reuse the active item without another item/started");
};
assert_eq!(params.item_id, "cmd-1");
assert_eq!(params.approval_id.as_deref(), Some("retry-1"));
assert_eq!(params.reason.as_deref(), Some("sandbox denied"));
assert_eq!(
params.available_decisions,
Some(vec![
CommandExecutionApprovalDecision::Accept,
CommandExecutionApprovalDecision::Cancel,
])
);
let completion_item = command_execution_completion_item("printf hi");
context
.raw_outgoing
.notify_client_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Decline,
})?,
)
.await;
let completed = recv_broadcast_message(&mut context.rx).await?;
assert!(matches!(
completed,
OutgoingMessage::AppServerNotification(ServerNotification::ItemCompleted(payload))
if matches!(payload.item, ThreadItem::CommandExecution { ref id, status: CommandExecutionStatus::Declined, .. } if id == "cmd-1")
));
Ok(())
}
let first_start = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(first_start);
#[tokio::test]
async fn repeated_approval_response_lifecycle_truth_table() -> Result<()> {
let mut context = CommandApprovalTestContext::new().await?;
let cases = [
(false, true, CommandExecutionApprovalDecision::Decline, true),
(true, true, CommandExecutionApprovalDecision::Decline, true),
(true, true, CommandExecutionApprovalDecision::Cancel, true),
(true, true, CommandExecutionApprovalDecision::Accept, false),
(
true,
false,
CommandExecutionApprovalDecision::Decline,
false,
),
(true, false, CommandExecutionApprovalDecision::Cancel, true),
];
let msg = recv_broadcast_message(&mut rx).await?;
match msg {
OutgoingMessage::AppServerNotification(ServerNotification::ItemStarted(payload)) => {
assert_eq!(payload.thread_id, conversation_id.to_string());
assert_eq!(payload.turn_id, "turn-1");
assert_eq!(
payload.item,
ThreadItem::CommandExecution {
id: "cmd-1".to_string(),
command: completion_item.command.clone(),
cwd: completion_item.cwd.clone(),
process_id: None,
source: CommandExecutionSource::Agent,
status: CommandExecutionStatus::InProgress,
command_actions: completion_item.command_actions.clone(),
aggregated_output: None,
exit_code: None,
duration_ms: None,
}
for (index, (has_active_item, is_sandbox_retry, decision, completes_on_response)) in
cases.into_iter().enumerate()
{
let item_id = format!("cmd-{index}");
let completion_item = command_execution_completion_item("printf hi");
if has_active_item {
assert!(
start_command_execution_item(
&context.conversation_id,
"turn-1".to_string(),
item_id.clone(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&context.outgoing,
&context.thread_state,
)
.await
);
let _initial_started = recv_broadcast_message(&mut context.rx).await?;
}
let started_now = start_command_execution_item(
&context.conversation_id,
"turn-1".to_string(),
item_id.clone(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&context.outgoing,
&context.thread_state,
)
.await;
assert_eq!(started_now, !has_active_item);
if started_now {
let _callback_started = recv_broadcast_message(&mut context.rx).await?;
}
let (response_tx, response_rx) = oneshot::channel();
response_tx
.send(Ok(serde_json::to_value(
CommandExecutionRequestApprovalResponse { decision },
)?))
.expect("response receiver should remain open");
let permission_guard = context
.thread_watch_manager
.note_permission_requested(&context.conversation_id.to_string())
.await;
on_command_execution_request_approval_response(
"turn-1".to_string(),
context.conversation_id,
Some(format!("callback-{index}")),
item_id.clone(),
Some(completion_item),
started_now,
is_sandbox_retry,
RequestId::Integer(index as i64),
response_rx,
Arc::clone(&context.conversation),
context.outgoing.clone(),
Arc::clone(&context.thread_state),
permission_guard,
)
.await;
if completes_on_response {
let completed = recv_broadcast_message(&mut context.rx).await?;
assert!(matches!(
completed,
OutgoingMessage::AppServerNotification(ServerNotification::ItemCompleted(payload))
if matches!(payload.item, ThreadItem::CommandExecution { ref id, status: CommandExecutionStatus::Declined, .. } if id == &item_id)
));
assert!(
!remove_started_command_execution_item(&context.thread_state, &item_id).await,
"a late runtime end must not complete {item_id} twice"
);
} else {
assert!(context.rx.try_recv().is_err());
assert!(
remove_started_command_execution_item(&context.thread_state, &item_id).await,
"the active parent retains completion ownership for {item_id}"
);
}
other => bail!("unexpected message: {other:?}"),
}
let second_start = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(!second_start);
assert!(rx.try_recv().is_err(), "duplicate start should not emit");
Ok(())
}
#[tokio::test]
async fn active_item_retry_decline_completes_but_execve_decline_is_suppressed() -> Result<()> {
let conversation_id = ThreadId::new();
let thread_state = new_thread_state();
let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY);
let outgoing = Arc::new(OutgoingMessageSender::new(
tx,
codex_analytics::AnalyticsEventsClient::disabled(),
));
let outgoing = ThreadScopedOutgoingMessageSender::new(
outgoing,
vec![ConnectionId(1)],
ThreadId::new(),
);
let completion_item = command_execution_completion_item("printf hi");
let initial_started = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(initial_started);
let _initial_started = recv_broadcast_message(&mut rx).await?;
let retry_started = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(!retry_started);
assert!(!suppress_repeated_prompt_completion_item(
Some("retry-approval-id"),
retry_started,
/*is_sandbox_retry*/ true,
/*cancel_requested*/ false,
));
assert!(!suppress_repeated_prompt_completion_item(
Some("retry-cancel-approval-id"),
retry_started,
/*is_sandbox_retry*/ true,
/*cancel_requested*/ true,
));
complete_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
/*process_id*/ None,
CommandExecutionSource::Agent,
completion_item.command_actions.clone(),
CommandExecutionStatus::Declined,
&outgoing,
&thread_state,
)
.await;
let retry_completed = recv_broadcast_message(&mut rx).await?;
match retry_completed {
OutgoingMessage::AppServerNotification(ServerNotification::ItemCompleted(payload)) => {
let ThreadItem::CommandExecution { id, status, .. } = payload.item else {
bail!("expected command execution completion");
};
assert_eq!(id, "cmd-1");
assert_eq!(status, CommandExecutionStatus::Declined);
}
other => bail!("unexpected message: {other:?}"),
}
let active_started = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-2".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(active_started);
let _active_started = recv_broadcast_message(&mut rx).await?;
let repeated_started = start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-2".to_string(),
completion_item.command,
completion_item.cwd,
completion_item.command_actions,
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
assert!(!repeated_started);
assert!(suppress_repeated_prompt_completion_item(
Some("zsh-subcommand-approval-id"),
repeated_started,
/*is_sandbox_retry*/ false,
/*cancel_requested*/ false,
));
assert!(!suppress_repeated_prompt_completion_item(
Some("zsh-subcommand-cancel-id"),
repeated_started,
/*is_sandbox_retry*/ false,
/*cancel_requested*/ true,
));
Ok(())
}
#[tokio::test]
async fn response_completion_suppresses_late_exec_command_end() -> Result<()> {
let conversation_id = ThreadId::new();
let thread_state = new_thread_state();
let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY);
let outgoing = Arc::new(OutgoingMessageSender::new(
tx,
codex_analytics::AnalyticsEventsClient::disabled(),
));
let outgoing = ThreadScopedOutgoingMessageSender::new(
outgoing,
vec![ConnectionId(1)],
ThreadId::new(),
);
let completion_item = command_execution_completion_item("printf hi");
start_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
completion_item.command_actions.clone(),
CommandExecutionSource::Agent,
&outgoing,
&thread_state,
)
.await;
let _started = recv_broadcast_message(&mut rx).await?;
complete_command_execution_item(
&conversation_id,
"turn-1".to_string(),
"cmd-1".to_string(),
completion_item.command.clone(),
completion_item.cwd.clone(),
/*process_id*/ None,
CommandExecutionSource::Agent,
completion_item.command_actions.clone(),
CommandExecutionStatus::Declined,
&outgoing,
&thread_state,
)
.await;
let completed = recv_broadcast_message(&mut rx).await?;
match completed {
OutgoingMessage::AppServerNotification(ServerNotification::ItemCompleted(payload)) => {
let ThreadItem::CommandExecution { id, status, .. } = payload.item else {
bail!("expected command execution completion");
};
assert_eq!(id, "cmd-1");
assert_eq!(status, CommandExecutionStatus::Declined);
}
other => bail!("unexpected message: {other:?}"),
}
let late_end_should_emit =
remove_started_command_execution_item(&thread_state, "cmd-1").await;
assert!(!late_end_should_emit);
assert!(
rx.try_recv().is_err(),
"a late ExecCommandEnd should not emit after response completion"
);
Ok(())
}