Files
codex/codex-rs/app-server/tests/suite/v2/thread_queue.rs
2026-05-30 15:35:54 -07:00

888 lines
29 KiB
Rust

use anyhow::Result;
use app_test_support::McpProcess;
use app_test_support::create_final_assistant_message_sse_response;
use app_test_support::create_mock_responses_server_sequence;
use app_test_support::create_mock_responses_server_sequence_unchecked;
use app_test_support::create_shell_command_sse_response;
use app_test_support::to_response;
use codex_app_server_protocol::AdditionalContextEntry;
use codex_app_server_protocol::AdditionalContextKind;
use codex_app_server_protocol::ClientInfo;
use codex_app_server_protocol::CommandExecutionApprovalDecision;
use codex_app_server_protocol::CommandExecutionRequestApprovalResponse;
use codex_app_server_protocol::ExperimentalFeatureEnablementSetParams;
use codex_app_server_protocol::ExperimentalFeatureEnablementSetResponse;
use codex_app_server_protocol::InitializeCapabilities;
use codex_app_server_protocol::JSONRPCError;
use codex_app_server_protocol::JSONRPCResponse;
use codex_app_server_protocol::QueuedTurnStatus;
use codex_app_server_protocol::RequestId;
use codex_app_server_protocol::ServerRequest;
use codex_app_server_protocol::ThreadQueueAddParams;
use codex_app_server_protocol::ThreadQueueAddResponse;
use codex_app_server_protocol::ThreadQueueChangedNotification;
use codex_app_server_protocol::ThreadQueueDeleteParams;
use codex_app_server_protocol::ThreadQueueDeleteResponse;
use codex_app_server_protocol::ThreadQueueListParams;
use codex_app_server_protocol::ThreadQueueListResponse;
use codex_app_server_protocol::ThreadQueueReorderParams;
use codex_app_server_protocol::ThreadQueueReorderResponse;
use codex_app_server_protocol::ThreadStartParams;
use codex_app_server_protocol::ThreadStartResponse;
use codex_app_server_protocol::TurnStartParams;
use codex_app_server_protocol::TurnSubmission;
use codex_app_server_protocol::UserInput as V2UserInput;
use std::collections::BTreeMap;
use std::collections::HashMap;
use tempfile::TempDir;
use tokio::time::timeout;
const DEFAULT_READ_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
const INVALID_REQUEST_ERROR_CODE: i64 = -32600;
#[tokio::test]
async fn idle_queue_add_dispatches_serialized_turn_and_drains_visible_queue() -> Result<()> {
let responses = vec![
create_final_assistant_message_sse_response("queued done")?,
create_final_assistant_message_sse_response("second queued done")?,
];
let server = create_mock_responses_server_sequence(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "never")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let add_request_id = mcp
.send_raw_request(
"thread/queue/add",
Some(serde_json::to_value(ThreadQueueAddParams {
thread_id: thread.id.clone(),
submission: text_submission("queued serialized input").into(),
})?),
)
.await?;
let add_response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(add_request_id)),
)
.await??;
let ThreadQueueAddResponse { queued_turn } = to_response(add_response)?;
assert!(matches!(queued_turn.status, QueuedTurnStatus::Pending));
let add_notification = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("thread/queue/changed"),
)
.await??;
let add_notification: ThreadQueueChangedNotification = serde_json::from_value(
add_notification
.params
.expect("thread/queue/changed params"),
)?;
assert_eq!(add_notification.thread_id, thread.id);
assert_eq!(add_notification.queued_turns, vec![queued_turn.clone()]);
assert_eq!(add_notification.dispatching_queued_turn_id, None);
let drain_notification = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("thread/queue/changed"),
)
.await??;
let drain_notification: ThreadQueueChangedNotification = serde_json::from_value(
drain_notification
.params
.expect("thread/queue/changed params"),
)?;
assert_eq!(drain_notification.thread_id, thread.id);
assert!(drain_notification.queued_turns.is_empty());
assert_eq!(
drain_notification.dispatching_queued_turn_id,
Some(queued_turn.id)
);
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
let list_request_id = mcp
.send_raw_request(
"thread/queue/list",
Some(serde_json::to_value(ThreadQueueListParams {
thread_id: thread.id.clone(),
cursor: None,
limit: None,
})?),
)
.await?;
let list_response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(list_request_id)),
)
.await??;
let ThreadQueueListResponse { data, next_cursor } = to_response(list_response)?;
assert!(data.is_empty());
assert_eq!(next_cursor, None);
queue_turn(&mut mcp, &thread.id, "second queued serialized input").await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
assert!(list_queue_ids(&mut mcp, &thread.id).await?.is_empty());
let requests = server
.received_requests()
.await
.expect("failed to fetch received requests");
assert_eq!(requests.len(), 2);
assert!(
String::from_utf8_lossy(&requests[0].body).contains("queued serialized input"),
"queued turn payload should reach the model request after state round-trip"
);
assert!(
String::from_utf8_lossy(&requests[1].body).contains("second queued serialized input"),
"a later queued turn should still drain after a fast terminal dispatch"
);
Ok(())
}
#[tokio::test]
async fn queue_add_rejects_ephemeral_threads() -> Result<()> {
let server = create_mock_responses_server_sequence_unchecked(vec![
create_final_assistant_message_sse_response("unused")?,
])
.await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "never")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let start_request_id = mcp
.send_thread_start_request(ThreadStartParams {
model: Some("mock-model".to_string()),
ephemeral: Some(true),
..Default::default()
})
.await?;
let start_response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(start_request_id)),
)
.await??;
let ThreadStartResponse { thread, .. } = to_response(start_response)?;
let add_request_id = mcp
.send_raw_request(
"thread/queue/add",
Some(serde_json::to_value(ThreadQueueAddParams {
thread_id: thread.id.clone(),
submission: text_submission("ephemeral queued turn").into(),
})?),
)
.await?;
let add_error: JSONRPCError = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_error_message(RequestId::Integer(add_request_id)),
)
.await??;
assert_eq!(add_error.error.code, INVALID_REQUEST_ERROR_CODE);
assert_eq!(
add_error.error.message,
format!(
"ephemeral thread does not support queued turns: {}",
thread.id
)
);
Ok(())
}
#[tokio::test]
async fn queue_add_rejects_oversized_serialized_submission() -> Result<()> {
let server = create_mock_responses_server_sequence_unchecked(vec![
create_final_assistant_message_sse_response("unused")?,
])
.await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "never")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let mut submission = text_submission("small prompt");
submission.additional_context = Some(HashMap::from([(
"oversized".to_string(),
AdditionalContextEntry {
value: "x".repeat((1 << 20) + 1),
kind: AdditionalContextKind::Application,
},
)]));
let add_request_id = mcp
.send_raw_request(
"thread/queue/add",
Some(serde_json::to_value(ThreadQueueAddParams {
thread_id: thread.id,
submission: submission.into(),
})?),
)
.await?;
let add_error: JSONRPCError = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_error_message(RequestId::Integer(add_request_id)),
)
.await??;
assert_eq!(
add_error.error.message,
"Queued turn submission exceeds the maximum length of 1048576 bytes."
);
Ok(())
}
#[tokio::test]
async fn queue_add_rejects_requests_when_feature_is_disabled() -> Result<()> {
let server = create_mock_responses_server_sequence_unchecked(vec![
create_final_assistant_message_sse_response("unused")?,
])
.await;
let codex_home = TempDir::new()?;
write_queue_test_config_without_feature(codex_home.path(), &server.uri(), "never")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let add_request_id = mcp
.send_raw_request(
"thread/queue/add",
Some(serde_json::to_value(ThreadQueueAddParams {
thread_id: thread.id,
submission: text_submission("disabled queued turn").into(),
})?),
)
.await?;
let add_error: JSONRPCError = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_error_message(RequestId::Integer(add_request_id)),
)
.await??;
assert_eq!(add_error.error.code, INVALID_REQUEST_ERROR_CODE);
assert_eq!(
add_error.error.message,
"app-server queue feature is disabled"
);
Ok(())
}
#[tokio::test]
async fn runtime_feature_enablement_controls_queue_access_without_deleting_rows() -> Result<()> {
let responses = vec![
create_shell_command_sse_response(
vec![
"python3".to_string(),
"-c".to_string(),
"print(42)".to_string(),
],
/*workdir*/ None,
Some(5000),
"queue-feature-blocker",
)?,
create_final_assistant_message_sse_response("active turn done")?,
create_final_assistant_message_sse_response("queued turn done")?,
];
let server = create_mock_responses_server_sequence_unchecked(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config_without_feature(codex_home.path(), &server.uri(), "untrusted")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let active_turn_request_id = mcp
.send_turn_start_request(text_turn(&thread.id, "keep the thread running"))
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(active_turn_request_id)),
)
.await??;
let approval_request = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_request_message(),
)
.await??;
let ServerRequest::CommandExecutionRequestApproval { request_id, .. } = approval_request else {
panic!("expected command approval to keep the active turn open");
};
set_queue_feature(&mut mcp, /*enabled*/ true).await?;
let queued_turn_id = queue_turn(&mut mcp, &thread.id, "durable queued turn").await?;
set_queue_feature(&mut mcp, /*enabled*/ false).await?;
let list_request_id = mcp
.send_raw_request(
"thread/queue/list",
Some(serde_json::to_value(ThreadQueueListParams {
thread_id: thread.id.clone(),
cursor: None,
limit: None,
})?),
)
.await?;
let list_error: JSONRPCError = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_error_message(RequestId::Integer(list_request_id)),
)
.await??;
assert_eq!(list_error.error.code, INVALID_REQUEST_ERROR_CODE);
assert_eq!(
list_error.error.message,
"app-server queue feature is disabled"
);
set_queue_feature(&mut mcp, /*enabled*/ true).await?;
assert_eq!(
list_queue_ids(&mut mcp, &thread.id).await?,
vec![queued_turn_id]
);
mcp.send_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
Ok(())
}
#[tokio::test]
async fn busy_thread_queue_rows_support_list_reorder_and_delete_before_drain() -> Result<()> {
let responses = vec![
create_shell_command_sse_response(
vec![
"python3".to_string(),
"-c".to_string(),
"print(42)".to_string(),
],
/*workdir*/ None,
Some(5000),
"queue-blocker",
)?,
create_final_assistant_message_sse_response("active turn done")?,
];
let server = create_mock_responses_server_sequence_unchecked(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "untrusted")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let active_turn_request_id = mcp
.send_turn_start_request(text_turn(&thread.id, "keep the thread running"))
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(active_turn_request_id)),
)
.await??;
let approval_request = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_request_message(),
)
.await??;
let ServerRequest::CommandExecutionRequestApproval { request_id, .. } = approval_request else {
panic!("expected command approval to keep the active turn open");
};
let first = queue_turn(&mut mcp, &thread.id, "first queued").await?;
let second = queue_turn(&mut mcp, &thread.id, "second queued").await?;
assert_eq!(
list_queue_ids(&mut mcp, &thread.id).await?,
vec![first.clone(), second.clone()]
);
let first_page = list_queue_page(&mut mcp, &thread.id, /*cursor*/ None, Some(1)).await?;
assert_eq!(
first_page
.data
.into_iter()
.map(|queued_turn| queued_turn.id)
.collect::<Vec<_>>(),
vec![first.clone()]
);
let second_page =
list_queue_page(&mut mcp, &thread.id, first_page.next_cursor, Some(1)).await?;
assert_eq!(
second_page
.data
.into_iter()
.map(|queued_turn| queued_turn.id)
.collect::<Vec<_>>(),
vec![second.clone()]
);
assert_eq!(second_page.next_cursor, None);
let reorder_request_id = mcp
.send_raw_request(
"thread/queue/reorder",
Some(serde_json::to_value(ThreadQueueReorderParams {
thread_id: thread.id.clone(),
queued_turn_ids: vec![second.clone(), first.clone()],
})?),
)
.await?;
let reorder_response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(reorder_request_id)),
)
.await??;
let ThreadQueueReorderResponse { queued_turns } = to_response(reorder_response)?;
assert_eq!(
queued_turns
.into_iter()
.map(|queued_turn| queued_turn.id)
.collect::<Vec<_>>(),
vec![second.clone(), first.clone()]
);
delete_queue_turn(&mut mcp, &thread.id, &second).await?;
delete_queue_turn(&mut mcp, &thread.id, &first).await?;
assert!(list_queue_ids(&mut mcp, &thread.id).await?.is_empty());
mcp.send_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
Ok(())
}
#[tokio::test]
async fn queued_turns_stay_serial_after_the_first_dispatch_starts() -> Result<()> {
let responses = vec![
create_shell_command_sse_response(
vec![
"python3".to_string(),
"-c".to_string(),
"print(42)".to_string(),
],
/*workdir*/ None,
Some(5000),
"queued-serial-blocker",
)?,
create_final_assistant_message_sse_response("first queued turn done")?,
create_final_assistant_message_sse_response("second queued turn done")?,
];
let server = create_mock_responses_server_sequence_unchecked(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "untrusted")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
queue_turn(&mut mcp, &thread.id, "first queued").await?;
let approval_request = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_request_message(),
)
.await??;
let ServerRequest::CommandExecutionRequestApproval { request_id, .. } = approval_request else {
panic!("expected queued turn approval request to keep the first dispatch active");
};
let second = queue_turn(&mut mcp, &thread.id, "second queued").await?;
assert_eq!(list_queue_ids(&mut mcp, &thread.id).await?, vec![second]);
mcp.send_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
assert!(list_queue_ids(&mut mcp, &thread.id).await?.is_empty());
let requests = server
.received_requests()
.await
.expect("failed to fetch received requests");
assert_eq!(requests.len(), 3);
assert!(
String::from_utf8_lossy(&requests[2].body).contains("second queued"),
"second queued follow-up should become its own later model request"
);
Ok(())
}
#[tokio::test]
async fn queued_turns_wait_for_a_just_accepted_direct_turn_to_become_visible() -> Result<()> {
let responses = vec![
create_shell_command_sse_response(
vec![
"python3".to_string(),
"-c".to_string(),
"print(42)".to_string(),
],
/*workdir*/ None,
Some(5000),
"direct-turn-blocker",
)?,
create_final_assistant_message_sse_response("direct turn done")?,
create_final_assistant_message_sse_response("queued follow-up done")?,
];
let server = create_mock_responses_server_sequence_unchecked(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "untrusted")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let direct_turn_request_id = mcp
.send_turn_start_request(text_turn(&thread.id, "direct turn first"))
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(direct_turn_request_id)),
)
.await??;
let queued_turn_id = queue_turn(&mut mcp, &thread.id, "queued turn after direct").await?;
assert_eq!(
list_queue_ids(&mut mcp, &thread.id).await?,
vec![queued_turn_id]
);
let approval_request = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_request_message(),
)
.await??;
let ServerRequest::CommandExecutionRequestApproval { request_id, .. } = approval_request else {
panic!("expected direct turn approval request to keep the direct turn open");
};
mcp.send_response(
request_id,
serde_json::to_value(CommandExecutionRequestApprovalResponse {
decision: CommandExecutionApprovalDecision::Accept,
})?,
)
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
assert!(list_queue_ids(&mut mcp, &thread.id).await?.is_empty());
let requests = server
.received_requests()
.await
.expect("failed to fetch received requests");
assert_eq!(requests.len(), 3);
assert!(
String::from_utf8_lossy(&requests[2].body).contains("queued turn after direct"),
"queued follow-up should become its own later model request"
);
Ok(())
}
#[tokio::test]
async fn queued_turns_drain_after_a_direct_turn_has_already_completed() -> Result<()> {
let responses = vec![
create_final_assistant_message_sse_response("direct turn done")?,
create_final_assistant_message_sse_response("queued follow-up done")?,
];
let server = create_mock_responses_server_sequence_unchecked(responses).await;
let codex_home = TempDir::new()?;
write_queue_test_config(codex_home.path(), &server.uri(), "never")?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
initialize_experimental(&mut mcp).await?;
let thread = start_thread(&mut mcp).await?;
let direct_turn_request_id = mcp
.send_turn_start_request(text_turn(&thread.id, "direct turn first"))
.await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(direct_turn_request_id)),
)
.await??;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
queue_turn(&mut mcp, &thread.id, "queued turn after completion").await?;
timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_notification_message("turn/completed"),
)
.await??;
assert!(list_queue_ids(&mut mcp, &thread.id).await?.is_empty());
let requests = server
.received_requests()
.await
.expect("failed to fetch received requests");
assert_eq!(requests.len(), 2);
assert!(
String::from_utf8_lossy(&requests[1].body).contains("queued turn after completion"),
"queued follow-up should drain after an already completed direct turn"
);
Ok(())
}
async fn initialize_experimental(mcp: &mut McpProcess) -> Result<()> {
timeout(
DEFAULT_READ_TIMEOUT,
mcp.initialize_with_capabilities(
ClientInfo {
name: "thread-queue-tests".to_string(),
title: None,
version: "0.0.0".to_string(),
},
Some(InitializeCapabilities {
experimental_api: true,
opt_out_notification_methods: None,
request_attestation: false,
}),
),
)
.await??;
Ok(())
}
async fn start_thread(mcp: &mut McpProcess) -> Result<codex_app_server_protocol::Thread> {
let request_id = mcp
.send_thread_start_request(ThreadStartParams {
model: Some("mock-model".to_string()),
..Default::default()
})
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
let ThreadStartResponse { thread, .. } = to_response(response)?;
Ok(thread)
}
async fn queue_turn(mcp: &mut McpProcess, thread_id: &str, text: &str) -> Result<String> {
let request_id = mcp
.send_raw_request(
"thread/queue/add",
Some(serde_json::to_value(ThreadQueueAddParams {
thread_id: thread_id.to_string(),
submission: text_submission(text).into(),
})?),
)
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
let ThreadQueueAddResponse { queued_turn } = to_response(response)?;
Ok(queued_turn.id)
}
async fn list_queue_ids(mcp: &mut McpProcess, thread_id: &str) -> Result<Vec<String>> {
let ThreadQueueListResponse { data, .. } =
list_queue_page(mcp, thread_id, /*cursor*/ None, /*limit*/ None).await?;
Ok(data.into_iter().map(|queued_turn| queued_turn.id).collect())
}
async fn list_queue_page(
mcp: &mut McpProcess,
thread_id: &str,
cursor: Option<String>,
limit: Option<u32>,
) -> Result<ThreadQueueListResponse> {
let request_id = mcp
.send_raw_request(
"thread/queue/list",
Some(serde_json::to_value(ThreadQueueListParams {
thread_id: thread_id.to_string(),
cursor,
limit,
})?),
)
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
to_response(response)
}
async fn delete_queue_turn(
mcp: &mut McpProcess,
thread_id: &str,
queued_turn_id: &str,
) -> Result<()> {
let request_id = mcp
.send_raw_request(
"thread/queue/delete",
Some(serde_json::to_value(ThreadQueueDeleteParams {
thread_id: thread_id.to_string(),
queued_turn_id: queued_turn_id.to_string(),
})?),
)
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
let ThreadQueueDeleteResponse { deleted } = to_response(response)?;
assert!(deleted);
Ok(())
}
async fn set_queue_feature(mcp: &mut McpProcess, enabled: bool) -> Result<()> {
let request_id = mcp
.send_experimental_feature_enablement_set_request(ExperimentalFeatureEnablementSetParams {
enablement: BTreeMap::from([("app_server_queue".to_string(), enabled)]),
})
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_READ_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
let ExperimentalFeatureEnablementSetResponse { enablement } = to_response(response)?;
assert_eq!(
enablement,
BTreeMap::from([("app_server_queue".to_string(), enabled)])
);
Ok(())
}
fn text_submission(text: &str) -> TurnSubmission {
TurnSubmission {
input: vec![V2UserInput::Text {
text: text.to_string(),
text_elements: Vec::new(),
}],
..Default::default()
}
}
fn text_turn(thread_id: &str, text: &str) -> TurnStartParams {
TurnStartParams {
thread_id: thread_id.to_string(),
input: text_submission(text).input,
..Default::default()
}
}
fn write_queue_test_config_without_feature(
codex_home: &std::path::Path,
server_uri: &str,
approval_policy: &str,
) -> std::io::Result<()> {
write_queue_test_config_with_optional_feature(
codex_home,
server_uri,
approval_policy,
/*app_server_queue*/ None,
)
}
fn write_queue_test_config(
codex_home: &std::path::Path,
server_uri: &str,
approval_policy: &str,
) -> std::io::Result<()> {
write_queue_test_config_with_optional_feature(
codex_home,
server_uri,
approval_policy,
/*app_server_queue*/ Some(true),
)
}
fn write_queue_test_config_with_optional_feature(
codex_home: &std::path::Path,
server_uri: &str,
approval_policy: &str,
app_server_queue: Option<bool>,
) -> std::io::Result<()> {
let feature_config = app_server_queue
.map(|enabled| format!("\n[features]\napp_server_queue = {enabled}\n"))
.unwrap_or_default();
std::fs::write(
codex_home.join("config.toml"),
format!(
r#"
model = "mock-model"
approval_policy = "{approval_policy}"
sandbox_mode = "read-only"
model_provider = "mock_provider"
{feature_config}
[model_providers.mock_provider]
name = "Mock provider for test"
base_url = "{server_uri}/v1"
wire_api = "responses"
request_max_retries = 0
stream_max_retries = 0
"#
),
)
}