From 76dc5f72a4f8b76298be3ec0b56d6d7d7e43f885 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 16 Mar 2026 15:16:35 -0400 Subject: [PATCH] Add turn-level retry for stream-ended-without-Finish errors When an LLM stream drops mid-response (e.g. under high concurrency with OpenAI), retry the same turn up to 3 times instead of failing the entire agent session. Conversation history is preserved across retries. Previously this killed the whole stage and restarted from scratch. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/crates/fabro-agent/src/session.rs | 115 ++++++++++++++++---------- 1 file changed, 73 insertions(+), 42 deletions(-) diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index 234ed994c..2d7dd1212 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -602,7 +602,7 @@ impl Session { async move { c.stream(&r).await } }) .await; - let mut event_stream = match stream_result { + let event_stream = match stream_result { Ok(stream) => stream, Err(err) => { self.event_emitter.emit( @@ -618,60 +618,91 @@ impl Session { } }; - let mut accumulator = StreamAccumulator::new(); + // Consume the stream, retrying up to 3 times on transient + // "stream ended without Finish" errors (e.g. dropped connections). + const STREAM_CONSUME_RETRIES: usize = 3; + let mut response = None; + let mut last_stream = event_stream; - while let Some(event_result) = event_stream.next().await { - match event_result { - Ok(event) => { - match &event { - StreamEvent::TextDelta { ref delta, .. } => { - self.event_emitter.emit( - self.id.clone(), - AgentEvent::TextDelta { - delta: delta.clone(), - }, - ); + for stream_attempt in 0..=STREAM_CONSUME_RETRIES { + let mut accumulator = StreamAccumulator::new(); + + while let Some(event_result) = last_stream.next().await { + match event_result { + Ok(event) => { + match &event { + StreamEvent::TextDelta { ref delta, .. } => { + self.event_emitter.emit( + self.id.clone(), + AgentEvent::TextDelta { + delta: delta.clone(), + }, + ); + } + StreamEvent::ReasoningDelta { ref delta } => { + self.event_emitter.emit( + self.id.clone(), + AgentEvent::ReasoningDelta { + delta: delta.clone(), + }, + ); + } + _ => {} } - StreamEvent::ReasoningDelta { ref delta } => { - self.event_emitter.emit( - self.id.clone(), - AgentEvent::ReasoningDelta { - delta: delta.clone(), - }, - ); - } - _ => {} + accumulator.process(&event); + } + Err(err) => { + self.event_emitter.emit( + self.id.clone(), + AgentEvent::Error { + error: AgentError::Llm(err.clone()), + }, + ); + return Err(AgentError::Llm(err)); } - accumulator.process(&event); } - Err(err) => { - self.event_emitter.emit( - self.id.clone(), - AgentEvent::Error { - error: AgentError::Llm(err.clone()), - }, - ); - return Err(AgentError::Llm(err)); + + // Check cancellation between chunks + if self.cancel_token.is_cancelled() { + break; } } - // Check cancellation between chunks + // If aborted during streaming, drop the stream to cancel the HTTP + // connection, then close the session before returning. if self.cancel_token.is_cancelled() { + drop(last_stream); + self.close(); + return Err(self.aborted_error()); + } + + if let Some(resp) = accumulator.response().cloned() { + response = Some(resp); break; } + + // No Finish event — retry if we have attempts left + if stream_attempt < STREAM_CONSUME_RETRIES { + tracing::warn!( + attempt = stream_attempt + 1, + max = STREAM_CONSUME_RETRIES, + "Stream ended without Finish event, retrying turn" + ); + // Re-open the stream with the same request + match client.stream(&request).await { + Ok(new_stream) => { + last_stream = new_stream; + } + Err(err) => { + return Err(AgentError::Llm(err)); + } + } + } } - // If aborted during streaming, drop the stream to cancel the HTTP - // connection, then close the session before returning. - if self.cancel_token.is_cancelled() { - drop(event_stream); - self.close(); - return Err(self.aborted_error()); - } - - let response = accumulator.response().cloned().ok_or_else(|| { + let response = response.ok_or_else(|| { AgentError::Llm(SdkError::Stream { - message: "Stream ended without a Finish event".into(), + message: "Stream ended without a Finish event (after retries)".into(), }) })?;