mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
Remove response.processed websocket request (#26447)
## Why The Responses websocket client no longer needs to send a follow-up `response.processed` request after a turn response has already been recorded. Keeping that extra acknowledgement path adds feature-gated control flow and a second websocket request shape that no longer carries useful behavior. ## What Changed - Removed the `response.processed` websocket request type and sender. - Removed the `responses_websocket_response_processed` feature flag and schema entry. - Removed turn and remote-compaction plumbing that only tracked response IDs to send the acknowledgement. - Removed tests that existed solely to cover the deleted feature path. ## Validation - `just fix -p codex-core -p codex-api -p codex-features`
This commit is contained in:
committed by
GitHub
Unverified
parent
f97d5c3275
commit
d312a53e2a
@@ -100,7 +100,6 @@ use tokio::sync::oneshot::error::TryRecvError;
|
||||
use tokio_tungstenite::tungstenite::Error;
|
||||
use tokio_tungstenite::tungstenite::Message;
|
||||
use tokio_util::sync::CancellationToken;
|
||||
use tracing::debug;
|
||||
use tracing::instrument;
|
||||
use tracing::trace;
|
||||
use tracing::warn;
|
||||
@@ -967,18 +966,6 @@ impl ModelClientSession {
|
||||
.set_connection_reused(/*connection_reused*/ false);
|
||||
}
|
||||
|
||||
pub(crate) async fn send_response_processed(&self, response_id: &str) {
|
||||
let Some(connection) = self.websocket_session.connection.as_ref() else {
|
||||
return;
|
||||
};
|
||||
if let Err(err) = connection
|
||||
.send_response_processed(response_id.to_string())
|
||||
.await
|
||||
{
|
||||
debug!("failed to send response.processed websocket request: {err}");
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
/// Builds shared Responses API transport options and request-body options.
|
||||
///
|
||||
@@ -1672,9 +1659,7 @@ fn parse_turn_metadata_header(turn_metadata_header: Option<&str>) -> Option<Head
|
||||
/// Meant to be called just before sending the request over the socket, to capture realistic
|
||||
/// transport timing.
|
||||
fn stamp_ws_stream_request_start_ms(request: &mut ResponsesWsRequest) {
|
||||
let ResponsesWsRequest::ResponseCreate(payload) = request else {
|
||||
return;
|
||||
};
|
||||
let ResponsesWsRequest::ResponseCreate(payload) = request;
|
||||
payload
|
||||
.client_metadata
|
||||
.get_or_insert_with(HashMap::new)
|
||||
|
||||
@@ -26,7 +26,6 @@ use codex_analytics::CompactionImplementation;
|
||||
use codex_analytics::CompactionPhase;
|
||||
use codex_analytics::CompactionReason;
|
||||
use codex_analytics::CompactionTrigger;
|
||||
use codex_features::Feature;
|
||||
use codex_protocol::error::CodexErr;
|
||||
use codex_protocol::error::Result as CodexResult;
|
||||
use codex_protocol::items::ContextCompactionItem;
|
||||
@@ -281,7 +280,6 @@ async fn run_remote_compact_task_inner_impl(
|
||||
);
|
||||
let RemoteCompactionV2Output {
|
||||
compaction_output,
|
||||
response_id,
|
||||
token_usage,
|
||||
} = compaction_output_result?;
|
||||
if let Some(token_usage) = token_usage {
|
||||
@@ -314,18 +312,11 @@ async fn run_remote_compact_task_inner_impl(
|
||||
|
||||
sess.emit_turn_item_completed(turn_context, compaction_item)
|
||||
.await;
|
||||
if turn_context
|
||||
.features
|
||||
.enabled(Feature::ResponsesWebsocketResponseProcessed)
|
||||
{
|
||||
client_session.send_response_processed(&response_id).await;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
struct RemoteCompactionV2Output {
|
||||
compaction_output: ResponseItem,
|
||||
response_id: String,
|
||||
token_usage: Option<TokenUsage>,
|
||||
}
|
||||
|
||||
@@ -409,7 +400,7 @@ async fn collect_compaction_output(
|
||||
let mut output_item_count = 0usize;
|
||||
let mut compaction_count = 0usize;
|
||||
let mut compaction_output = None;
|
||||
let mut completed_response_id = None;
|
||||
let mut saw_completed = false;
|
||||
let mut completed_token_usage = None;
|
||||
while let Some(event) = stream.next().await {
|
||||
match event? {
|
||||
@@ -422,12 +413,8 @@ async fn collect_compaction_output(
|
||||
}
|
||||
}
|
||||
}
|
||||
ResponseEvent::Completed {
|
||||
response_id,
|
||||
token_usage,
|
||||
..
|
||||
} => {
|
||||
completed_response_id = Some(response_id);
|
||||
ResponseEvent::Completed { token_usage, .. } => {
|
||||
saw_completed = true;
|
||||
completed_token_usage = token_usage;
|
||||
break;
|
||||
}
|
||||
@@ -435,12 +422,12 @@ async fn collect_compaction_output(
|
||||
}
|
||||
}
|
||||
|
||||
let Some(response_id) = completed_response_id else {
|
||||
if !saw_completed {
|
||||
return Err(CodexErr::Stream(
|
||||
"remote compaction v2 stream closed before response.completed".to_string(),
|
||||
None,
|
||||
));
|
||||
};
|
||||
}
|
||||
|
||||
if compaction_count != 1 {
|
||||
return Err(CodexErr::Fatal(format!(
|
||||
@@ -453,7 +440,6 @@ async fn collect_compaction_output(
|
||||
};
|
||||
Ok(RemoteCompactionV2Output {
|
||||
compaction_output,
|
||||
response_id,
|
||||
token_usage: completed_token_usage,
|
||||
})
|
||||
}
|
||||
@@ -804,7 +790,6 @@ mod tests {
|
||||
.expect("compaction should be collected");
|
||||
|
||||
assert_eq!(output.compaction_output, compaction);
|
||||
assert_eq!(output.response_id, "resp-compact");
|
||||
assert_eq!(
|
||||
output.token_usage,
|
||||
Some(TokenUsage {
|
||||
|
||||
@@ -1827,7 +1827,6 @@ async fn try_run_sampling_request(
|
||||
!sess.services.extensions.turn_item_contributors().is_empty();
|
||||
let mut active_item_is_streaming_to_client = false;
|
||||
let receiving_span = trace_span!("receiving_stream");
|
||||
let mut completed_response_id: Option<String> = None;
|
||||
let outcome: CodexResult<SamplingRequestResult> = loop {
|
||||
let handle_responses = trace_span!(
|
||||
parent: &receiving_span,
|
||||
@@ -2071,9 +2070,9 @@ async fn try_run_sampling_request(
|
||||
sess.services.models_manager.refresh_if_new_etag(etag).await;
|
||||
}
|
||||
ResponseEvent::Completed {
|
||||
response_id,
|
||||
token_usage,
|
||||
end_turn,
|
||||
..
|
||||
} => {
|
||||
flush_assistant_text_segments_all(
|
||||
&sess,
|
||||
@@ -2089,7 +2088,6 @@ async fn try_run_sampling_request(
|
||||
if let Some(false) = end_turn {
|
||||
needs_follow_up = true;
|
||||
}
|
||||
completed_response_id = Some(response_id);
|
||||
break Ok(SamplingRequestResult {
|
||||
needs_follow_up,
|
||||
last_agent_message,
|
||||
@@ -2213,15 +2211,6 @@ async fn try_run_sampling_request(
|
||||
)
|
||||
.await;
|
||||
|
||||
if sess
|
||||
.features
|
||||
.enabled(Feature::ResponsesWebsocketResponseProcessed)
|
||||
&& outcome.is_ok()
|
||||
&& let Some(response_id) = completed_response_id.as_deref()
|
||||
{
|
||||
client_session.send_response_processed(response_id).await;
|
||||
}
|
||||
|
||||
drain_in_flight(&mut in_flight, sess.clone(), turn_context.clone()).await?;
|
||||
|
||||
if should_emit_token_count {
|
||||
|
||||
Reference in New Issue
Block a user