From d6cd6b3f749ff514c4cbe1736ab5b8cb14e80e8c Mon Sep 17 00:00:00 2001 From: Max Loffgren Date: Mon, 7 Sep 2026 21:38:15 -0500 Subject: [PATCH] fix(core): start an interrupted stream over when the body read fails The restart for a stream the network cut before any reply text only fired when the body ended cleanly with no completion signal. A read error on the body was skipped as a decode failure, but every error the response stream yields carries that label, and this client applies no content decoding of its own: a chunked or sized body cut short arrives exactly there. Servers frame streams chunked, so the fault the restart was written for still ended the turn. Every body read error now counts as an interruption and gets the same fresh attempt as a clean cut. The sequence fixture can frame a body chunked and close after the first chunk; the new test streams reasoning that way, is cut, and asserts the reply is started over on a second connection with the retry announced in between. It fails on the previous code with a single served connection and a truncated result. --- crates/core/src/client.rs | 50 +++++++++++++++++++++++++++++++++++---- 1 file changed, 45 insertions(+), 5 deletions(-) diff --git a/crates/core/src/client.rs b/crates/core/src/client.rs index ce9dc0f..72546b5 100644 --- a/crates/core/src/client.rs +++ b/crates/core/src/client.rs @@ -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; } }; @@ -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(), ); @@ -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.