Skip to content

Commit 1de0085

Browse files
aibrahim-oaicodex
andauthored
Stream Realtime V2 background agent progress (openai#17264)
Stream Realtime V2 background agent updates while the background agent task is still running, then send the final tool output when it completes. User input during an active V2 handoff is acknowledged back to realtime as a steering update. Stack: - Depends on openai#17278 for the background_agent rename. - Depends on openai#17280 for the input task handler refactor. Coverage: - Adds an app-server integration regression test that verifies V2 progress is sent before the final function-call output. Validation: - just fmt - cargo check -p codex-core - cargo check -p codex-app-server --tests - git diff --check --------- Co-authored-by: Codex <noreply@openai.com>
1 parent 4e910bf commit 1de0085

2 files changed

Lines changed: 198 additions & 153 deletions

File tree

codex-rs/app-server/tests/suite/v2/realtime_conversation.rs

Lines changed: 74 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1222,6 +1222,7 @@ async fn webrtc_v2_background_agent_tool_call_delegates_and_returns_function_out
12221222
v2_background_agent_tool_call("call_v2", "delegate from v2"),
12231223
],
12241224
vec![],
1225+
vec![],
12251226
])]),
12261227
)
12271228
.await?;
@@ -1249,8 +1250,53 @@ async fn webrtc_v2_background_agent_tool_call_delegates_and_returns_function_out
12491250
requests[0]
12501251
);
12511252

1252-
let tool_output = harness.sideband_outbound_request(/*request_index*/ 1).await;
1253+
let progress = harness.sideband_outbound_request(/*request_index*/ 1).await;
1254+
assert_v2_progress_update(&progress, "delegated from v2");
1255+
1256+
let tool_output = harness.sideband_outbound_request(/*request_index*/ 2).await;
12531257
assert_v2_function_call_output(&tool_output, "call_v2", "delegated from v2");
1258+
assert_eq!(
1259+
function_call_output_sideband_requests(&harness.realtime_server).len(),
1260+
1
1261+
);
1262+
1263+
harness.shutdown().await;
1264+
Ok(())
1265+
}
1266+
1267+
#[tokio::test]
1268+
async fn webrtc_v2_background_agent_progress_is_sent_before_function_output() -> Result<()> {
1269+
skip_if_no_network!(Ok(()));
1270+
1271+
let mut harness = RealtimeE2eHarness::new(
1272+
RealtimeTestVersion::V2,
1273+
main_loop_responses(vec![create_final_assistant_message_sse_response(
1274+
"progress before final",
1275+
)?]),
1276+
realtime_sideband(vec![realtime_sideband_connection(vec![
1277+
vec![
1278+
session_updated("sess_v2_progress_before_final"),
1279+
v2_background_agent_tool_call("call_progress_order", "stream progress"),
1280+
],
1281+
vec![],
1282+
vec![],
1283+
])]),
1284+
)
1285+
.await?;
1286+
1287+
let started = harness.start_webrtc_realtime("v=offer\r\n").await?;
1288+
assert_eq!(started.started.version, RealtimeConversationVersion::V2);
1289+
1290+
let turn_completed = harness
1291+
.read_notification::<TurnCompletedNotification>("turn/completed")
1292+
.await?;
1293+
assert_eq!(turn_completed.thread_id, harness.thread_id);
1294+
1295+
let progress = harness.sideband_outbound_request(/*request_index*/ 1).await;
1296+
assert_v2_progress_update(&progress, "progress before final");
1297+
1298+
let tool_output = harness.sideband_outbound_request(/*request_index*/ 2).await;
1299+
assert_v2_function_call_output(&tool_output, "call_progress_order", "progress before final");
12541300

12551301
harness.shutdown().await;
12561302
Ok(())
@@ -1278,6 +1324,7 @@ async fn webrtc_v2_tool_call_delegated_turn_can_execute_shell_tool() -> Result<(
12781324
v2_background_agent_tool_call("call_shell", "run shell through delegated turn"),
12791325
],
12801326
vec![],
1327+
vec![],
12811328
])]);
12821329

12831330
let mut harness = RealtimeE2eHarness::new_with_sandbox(
@@ -1329,7 +1376,10 @@ async fn webrtc_v2_tool_call_delegated_turn_can_execute_shell_tool() -> Result<(
13291376
requests[1]
13301377
);
13311378

1332-
let tool_output = harness.sideband_outbound_request(/*request_index*/ 1).await;
1379+
let progress = harness.sideband_outbound_request(/*request_index*/ 1).await;
1380+
assert_v2_progress_update(&progress, "shell tool finished");
1381+
1382+
let tool_output = harness.sideband_outbound_request(/*request_index*/ 2).await;
13331383
assert_v2_function_call_output(&tool_output, "call_shell", "shell tool finished");
13341384
assert_eq!(
13351385
function_call_output_sideband_requests(&harness.realtime_server).len(),
@@ -1379,6 +1429,7 @@ async fn webrtc_v2_tool_call_does_not_block_sideband_audio() -> Result<()> {
13791429
}),
13801430
],
13811431
vec![],
1432+
vec![],
13821433
])]),
13831434
)
13841435
.await?;
@@ -1405,7 +1456,10 @@ async fn webrtc_v2_tool_call_does_not_block_sideband_audio() -> Result<()> {
14051456
.await?;
14061457
assert_eq!(turn_completed.thread_id, harness.thread_id);
14071458

1408-
let tool_output = harness.sideband_outbound_request(/*request_index*/ 1).await;
1459+
let progress = harness.sideband_outbound_request(/*request_index*/ 1).await;
1460+
assert_v2_progress_update(&progress, "late delegated result");
1461+
1462+
let tool_output = harness.sideband_outbound_request(/*request_index*/ 2).await;
14091463
assert_v2_function_call_output(&tool_output, "call_audio", "late delegated result");
14101464

14111465
harness.shutdown().await;
@@ -1643,6 +1697,23 @@ fn assert_v2_function_call_output(request: &Value, call_id: &str, expected_outpu
16431697
);
16441698
}
16451699

1700+
fn assert_v2_progress_update(request: &Value, expected_text: &str) {
1701+
assert_eq!(
1702+
request,
1703+
&json!({
1704+
"type": "conversation.item.create",
1705+
"item": {
1706+
"type": "message",
1707+
"role": "user",
1708+
"content": [{
1709+
"type": "input_text",
1710+
"text": format!("{expected_text}\n\nUpdate from background agent (task hasn't finished yet):")
1711+
}]
1712+
}
1713+
})
1714+
);
1715+
}
1716+
16461717
fn assert_v1_session_update(request: &Value) -> Result<()> {
16471718
assert_eq!(request["type"].as_str(), Some("session.update"));
16481719
assert_eq!(request["session"]["type"].as_str(), Some("quicksilver"));

0 commit comments

Comments
 (0)