mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
Trace logical websocket request after untraced warmup (#23581)
## Why `prewarm_websocket` intentionally stays out of rollout inference tracing, but the next traced websocket request can still reuse the warmup `response_id` and send an empty `input` delta. If tracing records that wire payload verbatim, replay sees an incremental request whose parent was never traced and cannot reconstruct the conversation. This fixes that at the producer boundary instead of relaxing `rollout-trace` replay semantics around unresolved `previous_response_id` values. ## What - track whether the last websocket response came from an untraced warmup and clear that state when the websocket session is reset or reconnected - when a traced websocket request reuses that warmup parent, keep sending the compressed websocket request on the wire but record the logical `ResponsesApiRequest` in the rollout trace - add a regression test that proves replay reconstructs the logical user message even though the websocket follow-up carries `previous_response_id = warm-1` with empty `input` - update `InferenceTraceAttempt::record_started` docs to reflect that callers may record a logical request rather than the exact transport payload ## Testing - `cargo test -p codex-core --test all responses_websocket_request_prewarm_traces_logical_request`
This commit is contained in:
@@ -30,6 +30,11 @@ use codex_protocol::protocol::Op;
|
||||
use codex_protocol::protocol::SessionSource;
|
||||
use codex_protocol::protocol::W3cTraceContext;
|
||||
use codex_protocol::user_input::UserInput;
|
||||
use codex_rollout_trace::ConversationPart;
|
||||
use codex_rollout_trace::InferenceTraceContext;
|
||||
use codex_rollout_trace::RawTraceEventPayload;
|
||||
use codex_rollout_trace::TraceWriter;
|
||||
use codex_rollout_trace::replay_bundle;
|
||||
use core_test_support::load_default_config_for_test;
|
||||
use core_test_support::responses::WebSocketConnectionConfig;
|
||||
use core_test_support::responses::WebSocketTestServer;
|
||||
@@ -535,6 +540,112 @@ async fn responses_websocket_request_prewarm_reuses_connection() {
|
||||
server.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn responses_websocket_request_prewarm_traces_logical_request() {
|
||||
skip_if_no_network!();
|
||||
|
||||
let server = start_websocket_server(vec![vec![
|
||||
vec![ev_response_created("warm-1"), ev_completed("warm-1")],
|
||||
vec![ev_response_created("resp-1"), ev_completed("resp-1")],
|
||||
]])
|
||||
.await;
|
||||
|
||||
let harness = websocket_harness_with_options(&server, /*runtime_metrics_enabled*/ true).await;
|
||||
let mut client_session = harness.client.new_session();
|
||||
let prompt = prompt_with_input(vec![message_item("hello")]);
|
||||
|
||||
client_session
|
||||
.prewarm_websocket(
|
||||
&prompt,
|
||||
&harness.model_info,
|
||||
&harness.session_telemetry,
|
||||
harness.effort,
|
||||
harness.summary,
|
||||
/*service_tier*/ None,
|
||||
/*turn_metadata_header*/ None,
|
||||
)
|
||||
.await
|
||||
.expect("websocket prewarm failed");
|
||||
|
||||
let trace_dir = TempDir::new().expect("trace dir");
|
||||
let writer = Arc::new(
|
||||
TraceWriter::create(
|
||||
trace_dir.path(),
|
||||
"trace-1".to_string(),
|
||||
harness.session_id.to_string(),
|
||||
harness.thread_id.to_string(),
|
||||
)
|
||||
.expect("trace writer"),
|
||||
);
|
||||
writer
|
||||
.append(RawTraceEventPayload::ThreadStarted {
|
||||
thread_id: harness.thread_id.to_string(),
|
||||
agent_path: "/root".to_string(),
|
||||
metadata_payload: None,
|
||||
})
|
||||
.expect("thread started");
|
||||
writer
|
||||
.append(RawTraceEventPayload::CodexTurnStarted {
|
||||
codex_turn_id: "turn-1".to_string(),
|
||||
thread_id: harness.thread_id.to_string(),
|
||||
})
|
||||
.expect("turn started");
|
||||
|
||||
let inference_trace = InferenceTraceContext::enabled(
|
||||
writer,
|
||||
harness.thread_id.to_string(),
|
||||
"turn-1".to_string(),
|
||||
harness.model_info.slug.clone(),
|
||||
"test-provider".to_string(),
|
||||
);
|
||||
|
||||
let mut stream = client_session
|
||||
.stream(
|
||||
&prompt,
|
||||
&harness.model_info,
|
||||
&harness.session_telemetry,
|
||||
harness.effort,
|
||||
harness.summary,
|
||||
/*service_tier*/ None,
|
||||
/*turn_metadata_header*/ None,
|
||||
&inference_trace,
|
||||
)
|
||||
.await
|
||||
.expect("websocket stream failed");
|
||||
|
||||
while let Some(event) = stream.next().await {
|
||||
if matches!(event, Ok(ResponseEvent::Completed { .. })) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
let connection = server.single_connection();
|
||||
let follow_up = connection
|
||||
.get(1)
|
||||
.expect("missing follow-up request")
|
||||
.body_json();
|
||||
assert_eq!(follow_up["previous_response_id"].as_str(), Some("warm-1"));
|
||||
assert_eq!(follow_up["input"], serde_json::json!([]));
|
||||
|
||||
let rollout = replay_bundle(trace_dir.path()).expect("replay trace");
|
||||
let inference = rollout
|
||||
.inference_calls
|
||||
.values()
|
||||
.next()
|
||||
.expect("inference should be present");
|
||||
assert_eq!(inference.request_item_ids.len(), 1);
|
||||
assert_eq!(
|
||||
rollout.conversation_items[&inference.request_item_ids[0]]
|
||||
.body
|
||||
.parts,
|
||||
vec![ConversationPart::Text {
|
||||
text: "hello".to_string(),
|
||||
}],
|
||||
);
|
||||
|
||||
server.shutdown().await;
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn responses_websocket_reuses_connection_after_session_drop() {
|
||||
skip_if_no_network!();
|
||||
|
||||
Reference in New Issue
Block a user