Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 45 additions & 5 deletions crates/core/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -475,11 +475,12 @@ async fn read_sse(
Err(e) => {
saw_terminator = false;
finish_reason = TRUNCATED.into();
// A body that fails to decode will fail the same way again;
// only the connection is worth a fresh attempt.
if !e.is_decode() {
interrupted = Some(describe_transport(&e));
}
// Every error this stream yields is the connection failing
// under the body: a chunked or sized body cut short surfaces
// here, not as a clean end, and this client applies no content
// decoding that could fail on its own. So each one is worth a
// fresh attempt.
interrupted = Some(describe_transport(&e));
break;
}
};
Expand Down Expand Up @@ -1026,6 +1027,20 @@ mod tests {
if sse.is_empty() {
continue;
}
// A body prefixed with CHUNKED: is framed the way real servers
// frame a stream, `Transfer-Encoding: chunked`, and the socket
// closes after the first chunk with no terminating chunk: the
// cut arrives as a body error, not as a clean end.
if let Some(payload) = sse.strip_prefix("CHUNKED:") {
let _ = stream.write_all(
format!(
"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\n\r\n{:x}\r\n{payload}\r\n",
payload.len()
)
.as_bytes(),
);
continue;
}
let _ = stream.write_all(
format!("HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nConnection: close\r\n\r\n{sse}").as_bytes(),
);
Expand Down Expand Up @@ -1081,6 +1096,31 @@ mod tests {
);
}

/// The same interruption as a real server delivers it. Streams are framed
/// chunked, and a connection cut mid-body arrives as a read error rather
/// than a clean end. The bug this guards: that arm never restarted the
/// reply, so the fault it was written for still ended the turn.
#[tokio::test]
async fn a_stream_cut_mid_chunk_is_started_over() {
let (result, deltas, served) = stream_sequence(vec![
"CHUNKED:data: {\"choices\":[{\"delta\":{\"reasoning_content\":\"let me\"},\"finish_reason\":null}]}\n\n".into(),
concat!(
"data: {\"choices\":[{\"delta\":{\"content\":\"all of it\"},\"finish_reason\":null}]}\n\n",
"data: {\"choices\":[{\"delta\":{},\"finish_reason\":\"stop\"}]}\n\n",
"data: [DONE]\n\n",
)
.into(),
])
.await;
assert_eq!(result.finish_reason, "stop");
assert_eq!(result.content, "all of it");
assert_eq!(served, 2, "the cut stream was started over on a second connection");
assert_eq!(deltas.len(), 3, "{deltas:?}");
assert_eq!(deltas[0], "reasoning:let me");
assert!(deltas[1].starts_with(&format!("retry:2/{MAX_ATTEMPTS}:")), "{}", deltas[1]);
assert_eq!(deltas[2], "content:all of it");
}

/// A connection the far side drops before answering is a transport fault
/// after the connection existed: each one is resent after a backoff,
/// announced, and the turn goes on with the reply that arrives.
Expand Down
Loading