diff --git a/docs/config.md b/docs/config.md index 175707f1..3056ea09 100644 --- a/docs/config.md +++ b/docs/config.md @@ -594,7 +594,7 @@ above is the one exception with richer per-provider precedence; | Field | Default | What it bounds | |---|---|---| | `stream_chunk_secs` | 300 | Per-chunk read deadline for a streaming LLM response (fallback for the per-provider key above) | -| `tool_call_gap_secs` | 60 | Stall window while a tool call is mid-assembly in the stream. A timeout here is retried automatically (the partial, incomplete tool call is discarded and the request restarted); raise it only if your provider legitimately pauses longer than 60s between tool-call deltas. | +| `tool_call_gap_secs` | 60 | Stall window while a tool call is mid-assembly in the stream. A timeout here is retried automatically when no response text has been emitted yet (the partial tool call is discarded and the request restarted); if text was already shown, the partial is kept to avoid a duplicated response. Raise it only if your provider legitimately pauses longer than 60s between tool-call deltas. | | `mcp_call_secs` | 120 | Total budget for one MCP tool call, including reconnect + retry | | `mcp_init_secs` | 10 | MCP server `initialize` handshake | | `lsp_request_secs` | 30 | Any non-`initialize` LSP request | diff --git a/src/agent/agent_loop/retry.rs b/src/agent/agent_loop/retry.rs index 5ab55022..22e61363 100644 --- a/src/agent/agent_loop/retry.rs +++ b/src/agent/agent_loop/retry.rs @@ -62,11 +62,16 @@ use super::tool::AbortSignal; /// gives it a chance to break out instead of repeating the stall. const STALL_RECOVERY_NUDGE: &str = "Your previous response did not complete in time — you may have been stuck in a long reasoning loop. Stop deliberating: state your conclusion concisely and take the next concrete action (a tool call, or your final answer) now."; -/// Whether `msg` looks like a timeout (vs a generic transient network blip), -/// used to gate the stall-recovery nudge so it doesn't fire on, say, a 503. -fn is_timeout_error(msg: &str) -> bool { +/// Whether `msg` looks like a stall timeout the nudge can help break +/// out of (request-level timeout, or a mid-assembly stream-chunk +/// timeout where the model was actively producing), vs a generic +/// transient blip (503) — or the plain transport stream-chunk timeout +/// ("provider stalled or connection silently dropped"), which is a +/// connection drop, not a reasoning-loop stall: blaming "a long +/// reasoning loop" there is mis-attributed, so the retry re-runs clean. +fn is_stall_timeout(msg: &str) -> bool { let l = msg.to_ascii_lowercase(); - l.contains("timeout") || l.contains("timed out") + (l.contains("timeout") || l.contains("timed out")) && !l.contains("connection silently dropped") } /// Wrap an inner `StreamFn` with retry-on-transient-error @@ -133,19 +138,21 @@ pub fn retrying_stream_fn_with_non_retryable( return; } let kind = classify_error(error); - // #545: a stream-chunk TIMEOUT is - // retryable even after content has - // committed. A timeout means the - // provider stalled mid-emission, so the - // partial is incomplete; the consumer - // discards it on StreamEvent::Retry - // (stream.rs PROV-5 reset) and the next - // attempt starts clean. Other committed - // errors (hard resets, auth, …) still - // surface — retrying those would - // duplicate tokens already shown. - let retryable = policy.should_retry(attempts, kind) - && (!committed || is_timeout_error(error)); + // Retry only when NO content has + // committed. Once the consumer has seen + // real text/thinking deltas, re-running + // the request would duplicate them + // downstream: the retry reset in + // stream.rs (PROV-5) clears the message + // history but not text already rendered + // to the UI, so attempt 2 would re-emit + // the response a second time. The #545 + // mid-assembly-stall fix is preserved — + // a timeout before any text commits (e.g. + // a lone tool-call fragment) still retries + // below. + let retryable = + policy.should_retry(attempts, kind) && !committed; if retryable { retry_msg = Some(error.clone()); // Don't yield this Error — we're @@ -153,7 +160,7 @@ pub fn retrying_stream_fn_with_non_retryable( break; } // Retry exhausted, non-retryable kind, - // or a committed non-timeout → surface. + // or committed → surface. yield evt; return; } @@ -214,7 +221,7 @@ pub fn retrying_stream_fn_with_non_retryable( // If the previous attempt timed out, reinject a // one-time "conclude and act" nudge so the retry can // break out of a stall instead of repeating it. - let was_timeout = is_timeout_error(&err_msg); + let was_stall_timeout = is_stall_timeout(&err_msg); // PROV-2: surface the retry so the UI can // show a banner instead of freezing. yield StreamEvent::Retry { @@ -222,7 +229,7 @@ pub fn retrying_stream_fn_with_non_retryable( delay_ms: backoff.as_millis() as u64, error: err_msg, }; - if was_timeout && !stall_nudged { + if was_stall_timeout && !stall_nudged { stall_nudged = true; ctx.messages.push(serde_json::json!({ "role": "user", @@ -489,21 +496,23 @@ mod tests { assert!(matches!(events[2], StreamEvent::Error { .. })); } - /// #545: a stream-chunk TIMEOUT after content has committed is - /// still retried. Unlike a hard "connection reset", a timeout - /// means the provider stalled mid-emission — the partial is - /// incomplete and the consumer discards it on - /// `StreamEvent::Retry` (stream.rs PROV-5 reset), so retrying is - /// safe. This is what keeps a reasoning model that thinks, then - /// stalls mid-tool-call, from halting the whole run. + /// Regression for the duplicate-response bug: a stream-chunk + /// timeout AFTER committed text content must SURFACE the error, + /// not retry. #547 made timeouts retryable even when committed on + /// the assumption the consumer discards the partial on + /// `StreamEvent::Retry` — but that reset (stream.rs PROV-5) clears + /// only the message history, not text already rendered to the UI, + /// so attempt 2 re-emitted the response a second time. Gating + /// retry on `!committed` keeps the #545 mid-assembly-stall fix + /// (no text yet) while preventing duplication once text is shown. #[tokio::test] - async fn retries_timeout_after_content_committed() { + async fn surfaces_timeout_after_content_committed() { let (factory, counter) = counted_canned(vec![ vec![ StreamEvent::Start { partial: empty_assistant(), }, - // Real content streamed → committed. + // Real text streamed → committed → retry would dup. StreamEvent::Delta { partial: assistant_with("partial "), phase: DeltaPhase::TextStart, @@ -515,6 +524,60 @@ mod tests { .to_string(), }, ], + // Second attempt is supplied to prove it is NOT used. + vec![StreamEvent::Done { + reason: StopReason::Stop, + message: assistant_with("ok"), + usage: None, + }], + ]); + let wrapped = retrying_stream_fn(factory, RecoveryPolicy::default()); + + let task = tokio::spawn(async move { + drain(wrapped( + ctx(), + crate::agent::agent_loop::StreamOptions::from_signal(AbortSignal::new()), + )) + .await + }); + let events = task.await.unwrap(); + + // Not retried — inner stream called exactly once. + assert_eq!(counter.load(Ordering::SeqCst), 1); + // No Retry notice: the error was surfaced, not retried. + assert!( + !events + .iter() + .any(|e| matches!(e, StreamEvent::Retry { .. })), + "expected no retry after committed content, got: {events:?}" + ); + // The stream ends on the surfaced Error, not attempt 2's Done. + assert!(matches!(events.last(), Some(StreamEvent::Error { .. }))); + } + + /// #545 (preserved): a mid-assembly stream-chunk timeout BEFORE + /// any text content is committed still retries — the original + /// case of a reasoning model that emits a tool-call fragment then + /// stalls. No rendered text exists, so retrying is safe and + /// avoids halting the whole run. + #[tokio::test] + async fn retries_timeout_before_content_committed() { + let (factory, counter) = counted_canned(vec![ + vec![ + StreamEvent::Start { + partial: empty_assistant(), + }, + // A tool-call fragment is NOT committed content + // (is_content_delta) — only text/thinking phases are. + StreamEvent::Delta { + partial: empty_assistant(), + phase: DeltaPhase::ToolCallStart, + }, + StreamEvent::Error { + error: "stream chunk timed out after 29s while a tool call was mid-assembly" + .to_string(), + }, + ], vec![StreamEvent::Done { reason: StopReason::Stop, message: assistant_with("ok"), @@ -536,9 +599,6 @@ mod tests { // Retried once → second attempt succeeded. assert_eq!(counter.load(Ordering::SeqCst), 2); - // A Retry notice was emitted between attempts (the consumer - // uses it to reset partial state + show a banner instead of - // freezing). assert!( events .iter() @@ -784,4 +844,32 @@ mod tests { "no nudge for a non-timeout error" ); } + + /// A plain stream-chunk timeout — the transport case ("provider + /// stalled or connection silently dropped", no tool call open) — + /// retries, but is NOT a reasoning-loop stall, so the + /// stall-recovery nudge must not be injected. The nudge blames + /// "a long reasoning loop", which is mis-attributed for a + /// connection drop: the retry should just re-run cleanly. + #[tokio::test] + async fn plain_stream_chunk_timeout_does_not_nudge() { + // Verbatim shape produced by rig_stream.rs for a plain + // (no tool call open) stream-chunk timeout. + let (factory, seen) = recording_factory( + "stream chunk timed out after 300s (provider stalled or connection silently dropped) \ + — bump `stream_chunk_timeout_secs` in config.json if your model has long reasoning gaps", + ); + let wrapped = retrying_stream_fn(factory, RecoveryPolicy::default()); + let _ = drain(wrapped( + ctx(), + crate::agent::agent_loop::StreamOptions::from_signal(AbortSignal::new()), + )) + .await; + let recorded = seen.lock().unwrap(); + assert_eq!(recorded.len(), 2, "pre-commit timeout still retries"); + assert!( + !has_stall_nudge(&recorded[1]), + "transport stream-chunk timeout is not a reasoning-loop stall — no nudge" + ); + } } diff --git a/src/agent/agent_loop/rig_stream.rs b/src/agent/agent_loop/rig_stream.rs index 1a94b447..63bca3b7 100644 --- a/src/agent/agent_loop/rig_stream.rs +++ b/src/agent/agent_loop/rig_stream.rs @@ -187,7 +187,7 @@ where // retry. Matches runner.rs:301-304. let detail = if !open_tool_calls.is_empty() { format!( - "stream chunk timed out after {}s while a tool call was mid-assembly (provider stalled emitting tool-call deltas — common DeepSeek symptom; the harness narrows to {}s while assembling tool calls). This is retried automatically; if your model legitimately pauses longer than {}s between deltas, raise `timeouts.tool_call_gap_secs` in config.json.", + "stream chunk timed out after {}s while a tool call was mid-assembly (provider stalled emitting tool-call deltas — common DeepSeek symptom; the harness narrows to {}s while assembling tool calls). Retried automatically when no text has emitted yet; otherwise the partial response is kept to avoid duplicating it. If your model legitimately pauses longer than {}s between deltas, raise `timeouts.tool_call_gap_secs` in config.json.", t.as_secs(), tool_call_gap_timeout.as_secs(), tool_call_gap_timeout.as_secs(),