diff --git a/_plans/PREMATURE_COMPLETE_BUG.md b/_plans/PREMATURE_COMPLETE_BUG.md index 272dd8a1..98fc7eb2 100644 --- a/_plans/PREMATURE_COMPLETE_BUG.md +++ b/_plans/PREMATURE_COMPLETE_BUG.md @@ -531,3 +531,69 @@ Codex avoids this class when using Responses because completion is anchored on ` - `cargo check` - `cargo fmt --check` - `git diff --check` + +## 2026-08-31 Session Recurrence + +### User-Visible Symptom + +Session `bgwb0odvgy97qsr1joml2sc3` kept finishing with no error after the user said `continue`. The last visible assistant text was a preamble, then four `read`s of a local Expo project, then a reasoning-only step, then idle. + +### `app.log` Evidence + +Primary session id: `bgwb0odvgy97qsr1joml2sc3`. Model: `grok-4.6` via `https://cli-chat-proxy.grok.com` (`provider_kind=OpenAI`). + +Relevant sequence: + +- `02:29:50`: stream started (`turn_idx=5`, `input_messages=13`, `messages=179`, `agent_max_steps=None`). +- `02:30:03-02:30:06`: step 1 streamed unphased preamble text + four `read` tool calls. `assistant_message_phase=unknown`. +- `02:30:07`: `response.completed end_turn=None reasoning_items=1`. Tools executed successfully (`error_results=0`). +- `02:30:10`: step 2 started after tool results (`messages=189`). +- `02:30:13`: step 2 streamed reasoning only (`"Let me view the screenshots..."` in the persisted message). No text, no tool-call chunks. +- `02:30:14`: `response.completed end_turn=None reasoning_items=1`. Usage `output=137` (reasoning-sized; no leftover function-call budget). +- `02:30:14`: AISDK logged `provider_step_finish step=2 has_tool_call=false end_turn=None provider_finish_reason=unknown last_phase=unknown assistant_text_chars=0 action=finish preview=""`. +- `02:30:14`: crabcode marked the stream complete: `outcome=Exhausted`, `effective_outcome=Finished`, `stop_reason=Some(Finish)`. No Failed/Incomplete/Cancelled. + +An earlier continue turn in the same session (`tq36osvr1zxhibna28rqeb3i`) finished even earlier: reasoning + preamble text, zero tools. + +### Root Cause + +Same finish gate as the 2026-05-28 phase-less incident, but on the xAI Responses transport: + +1. `response.completed` is the terminal event. xAI/OpenAI Responses does not emit `ChunkType::End { reason }`, so `provider_finish_reason` stays `None` (logged `unknown`). +2. Message phases are also absent (`last_phase=unknown`). +3. `phase_less_ambiguous_requires_follow_up` only continues when `provider_finish_reason.is_some_and(|reason| !reason.is_final_assistant_stop())`. That path exists for Anthropic `end_turn`. For Responses, `None` fails `is_some_and`, so the step finishes. +4. Step 2 is stronger than the preamble case: tools were available, assistant text was empty, only reasoning arrived, `end_turn` was not `true`. AISDK still treated that as a real finish. + +Not a dropped-tool-call proof for this run: step 1 streamed function calls live, and step 2's `output=137` matches reasoning-only. `response.completed` still does not log/apply `output[].type` besides reasoning items, so the next recurrence should capture `output_types`. + +### Diagnostics Added + +- `src/aisdk/providers/openai.rs` + - Log `openai-responses completed status=... end_turn=... incomplete_reason=... output_count=... output_types=[...]`. +- `src/aisdk/response.rs` + - `provider_step_finish` now includes `reasoning_chars`, `tools`, and `follow_up[end_turn= commentary= phase_less= empty_output=]`. + +### Runtime Fix Applied 2026-08-31 (Grok Build empty resample) + +Reverted the agent-loop reminder / preamble-continue hacks. Grok Build does not keep the turn alive by inspecting assistant prose or injecting a "please continue" user message. + +Reference: `.devrefs/references/xai-org/grok-build/crates/codegen/xai-grok-sampler/src/actor/request_task.rs` + +- `ConversationResponse::empty_reason()` is `ReasoningOnly` or `NoVisibleContent` when there is no assistant text and no tool calls. +- That completed payload is `AttemptOutcome::Empty`: retry the **same sampling request**, do not accept it as a finished turn, do not append it to the conversation. +- Content-filter empties are not retried. +- After the retry budget, the request fails (`SamplingError::EmptyResponse`); it is not `Finish`. +- Reminders are only for doom-loop recovery. + +Crabcode now resamples the same provider step (rollback streamed reasoning) when a terminal response has no text and no tool calls. + +Preamble text with no tools (`I'll pull the screenshot language next…`) is still a completed assistant message, same as Grok Build: `empty_reason` is None when content is non-empty. That is prompt/model behavior, not a sampler retry. + +Validation: + +- `cargo test resamples_reasoning_only_completed_response_like_grok_build` +- `cargo test resamples_no_visible_content_completed_response` +- `cargo test empty_response_exhaustion_fails_instead_of_finish` +- `cargo test content_filter_empty_does_not_resample` +- `cargo test hosted_tool_only_completed_response_is_not_empty` +- `cargo test phase_less_text_without_finish_metadata_still_finishes` diff --git a/src/aisdk/chunk.rs b/src/aisdk/chunk.rs index f9cd8917..c062ccd7 100644 --- a/src/aisdk/chunk.rs +++ b/src/aisdk/chunk.rs @@ -19,6 +19,7 @@ pub enum ChunkType { ResponseCompleted { end_turn: Option, reasoning_items: Vec, + doom_loop_triggers: Vec, }, Retry(crate::retry::RetryStatus), StreamRollback { @@ -56,6 +57,7 @@ impl ChunkType { Self::ResponseCompleted { end_turn, reasoning_items: Vec::new(), + doom_loop_triggers: Vec::new(), } } } diff --git a/src/aisdk/providers/openai.rs b/src/aisdk/providers/openai.rs index 10237e1e..45e60fe1 100644 --- a/src/aisdk/providers/openai.rs +++ b/src/aisdk/providers/openai.rs @@ -1548,23 +1548,19 @@ fn response_sse_data_to_chunk(data: &str) -> Option> { if let Some(usage) = resp.get("usage") { log_openai_responses_usage(usage); } + log_openai_responses_completed(resp); Some(Ok(ChunkType::ResponseCompleted { end_turn: resp.get("end_turn").and_then(|value| value.as_bool()), reasoning_items: reasoning_items_from_response_output(resp), + doom_loop_triggers: doom_loop_triggers_from(resp), })) } // Grok Build / cli-chat-proxy: `response.doom_loop_check` with // `doom_loop_check.triggers` like `tail_repetition:8@thinking`. - // `.devrefs/references/xai-org/grok-build/crates/codegen/xai-grok-sampler/src/doom_loop.rs` + // Also present on `response.completed` (`doom_loop_check` field). + // `.devrefs/references/xai-org/grok-build/crates/codegen/xai-grok-sampling-types/src/doom_loop.rs` "response.doom_loop_check" => { - let triggers = value - .pointer("/doom_loop_check/triggers") - .and_then(|value| value.as_array()) - .into_iter() - .flatten() - .filter_map(|value| value.as_str()) - .collect::>() - .join(","); + let triggers = doom_loop_triggers_from(&value).join(","); Some(Ok(ChunkType::Metadata(format!( "doom_loop_check triggers={triggers}" )))) @@ -1629,6 +1625,68 @@ fn log_openai_responses_usage(usage: &serde_json::Value) { )); } +/// Attribute a `response.completed` payload: status, incomplete reason, and +/// output item types. Needed when a turn finishes with no error after +/// reasoning-only / empty assistant output. +fn log_openai_responses_completed(response: &serde_json::Value) { + let status = response + .get("status") + .and_then(|value| value.as_str()) + .unwrap_or("unknown"); + let end_turn = response.get("end_turn").and_then(|value| value.as_bool()); + let incomplete_reason = response + .get("incomplete_details") + .and_then(|details| details.get("reason")) + .and_then(|value| value.as_str()) + .unwrap_or("none"); + let (output_count, output_types) = summarize_response_output_types(response); + + let doom_loop_triggers = doom_loop_triggers_from(response); + let doom_loop = if doom_loop_triggers.is_empty() { + "none".to_string() + } else { + doom_loop_triggers.join(",") + }; + + crate::log::log(&format!( + "openai-responses completed status={status} end_turn={end_turn:?} incomplete_reason={incomplete_reason} output_count={output_count} output_types=[{output_types}] doom_loop_check={doom_loop}" + )); +} + +fn doom_loop_triggers_from(value: &serde_json::Value) -> Vec { + value + .pointer("/doom_loop_check/triggers") + .or_else(|| value.pointer("/response/doom_loop_check/triggers")) + .and_then(|value| value.as_array()) + .into_iter() + .flatten() + .filter_map(|value| value.as_str()) + .map(str::to_string) + .collect() +} + +fn summarize_response_output_types(response: &serde_json::Value) -> (usize, String) { + let Some(output) = response.get("output").and_then(|value| value.as_array()) else { + return (0, String::new()); + }; + + let mut counts: std::collections::BTreeMap<&str, usize> = std::collections::BTreeMap::new(); + for item in output { + let item_type = item + .get("type") + .and_then(|value| value.as_str()) + .unwrap_or("unknown"); + *counts.entry(item_type).or_default() += 1; + } + + let output_types = counts + .into_iter() + .map(|(item_type, count)| format!("{item_type}={count}")) + .collect::>() + .join(","); + (output.len(), output_types) +} + fn responses_provider_error_message(value: &serde_json::Value, fallback: &str) -> String { let code = response_error_field(value, "code"); let message = response_error_field(value, "message"); @@ -2348,10 +2406,10 @@ mod tests { add_responses_lite_header, build_openai_messages, build_websocket_request_body, fresh_websocket_request_body, is_client_tool_call_event, openai_chunk_is_terminal, request_snapshot_from_body, response_sse_data_to_chunk, responses_function_call_chunk, - websocket_connection_is_idle, websocket_continuation_mode_after_idle_policy, - websocket_continuation_mode_from_state, OpenAI, OpenAIResponseSnapshot, - OpenAIWebsocketState, WebsocketContinuationMode, WebsocketStreamProgress, - OPENAI_CODEX_WINDOW_ID_HEADER, OPENAI_RESPONSES_LITE_HEADER, + summarize_response_output_types, websocket_connection_is_idle, + websocket_continuation_mode_after_idle_policy, websocket_continuation_mode_from_state, + OpenAI, OpenAIResponseSnapshot, OpenAIWebsocketState, WebsocketContinuationMode, + WebsocketStreamProgress, OPENAI_CODEX_WINDOW_ID_HEADER, OPENAI_RESPONSES_LITE_HEADER, OPENAI_RESPONSES_LITE_WS_METADATA_KEY, OPENAI_WEBSOCKET_FAILURES_BEFORE_FALLBACK, OPENAI_WEBSOCKET_IDLE_MAX, }; @@ -2913,6 +2971,49 @@ mod tests { } } + #[test] + fn response_completed_captures_terminal_doom_loop_triggers() { + let chunk = response_sse_data_to_chunk( + r#"{"type":"response.completed","response":{"output":[{"type":"reasoning","id":"rs_1"}],"doom_loop_check":{"triggers":["tail_repetition:8@thinking"]}}}"#, + ) + .expect("expected terminal chunk"); + + match chunk { + Ok(ChunkType::ResponseCompleted { + doom_loop_triggers, .. + }) => { + assert_eq!( + doom_loop_triggers, + vec!["tail_repetition:8@thinking".to_string()] + ); + } + other => panic!("expected ResponseCompleted, got {other:?}"), + } + } + + #[test] + fn summarize_response_output_types_counts_reasoning_and_function_calls() { + let response = serde_json::json!({ + "output": [ + {"type": "reasoning", "id": "rs_1"}, + {"type": "function_call", "call_id": "call_1", "name": "read"}, + {"type": "function_call", "call_id": "call_2", "name": "read"}, + ] + }); + let (count, types) = summarize_response_output_types(&response); + assert_eq!(count, 3); + assert_eq!(types, "function_call=2,reasoning=1"); + } + + #[test] + fn summarize_response_output_types_empty_when_missing_output() { + let response = serde_json::json!({"status": "completed"}); + assert_eq!( + summarize_response_output_types(&response), + (0, String::new()) + ); + } + #[test] fn doom_loop_check_sse_becomes_metadata() { let chunk = response_sse_data_to_chunk( diff --git a/src/aisdk/response.rs b/src/aisdk/response.rs index 479a0204..cef392f7 100644 --- a/src/aisdk/response.rs +++ b/src/aisdk/response.rs @@ -221,266 +221,404 @@ pub async fn stream_with_tools( let mut accumulated_text = String::new(); let mut accumulated_reasoning = String::new(); let mut reasoning_replay_items: Vec = Vec::new(); - let mut server_doom_loop = false; + let mut server_doom_loop; let mut saw_terminal_event = false; let mut response_end_turn = None; let mut provider_finish_reason = None; let mut last_assistant_message_phase = None; let mut current_assistant_message_phase = None; let mut emitted_non_replayable_output = false; - - loop { - let next_chunk = if let Some(token) = cancel_token.as_ref() { - tokio::select! { - _ = token.cancelled() => { - let err = "Streaming cancelled by user".to_string(); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; + let mut had_provider_tool_call = false; + + 'resample: loop { + server_doom_loop = false; + loop { + let next_chunk = if let Some(token) = cancel_token.as_ref() { + tokio::select! { + _ = token.cancelled() => { + let err = "Streaming cancelled by user".to_string(); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + chunk = stream.next() => chunk, } - chunk = stream.next() => chunk, - } - } else { - stream.next().await - }; - let Some(chunk) = next_chunk else { break }; - match chunk { - Ok(ChunkType::AssistantMessagePhase { phase }) => { - current_assistant_message_phase = phase; - last_assistant_message_phase = phase; - let label = message_phase_label(phase); - let _ = tx_loop.send(ChunkType::Metadata(format!( - "assistant_message_phase={label}" - ))); - } - Ok(ChunkType::ResponseCompleted { - end_turn, - reasoning_items, - }) => { - saw_terminal_event = true; - response_end_turn = end_turn; - for item in reasoning_items { - merge_reasoning_replay_item(&mut reasoning_replay_items, item); + } else { + stream.next().await + }; + let Some(chunk) = next_chunk else { break }; + match chunk { + Ok(ChunkType::AssistantMessagePhase { phase }) => { + current_assistant_message_phase = phase; + last_assistant_message_phase = phase; + let label = message_phase_label(phase); + let _ = tx_loop.send(ChunkType::Metadata(format!( + "assistant_message_phase={label}" + ))); } - let _ = tx_loop.send(ChunkType::Metadata(format!( - "response.completed end_turn={end_turn:?} reasoning_items={}", - reasoning_replay_items.len() - ))); - } - Ok(ChunkType::ReasoningItem(item)) => { - let id = item.id.clone().unwrap_or_default(); - let encrypted_bytes = item - .encrypted_content - .as_ref() - .map(String::len) - .unwrap_or(0); - merge_reasoning_replay_item(&mut reasoning_replay_items, item); - let _ = tx_loop.send(ChunkType::Metadata(format!( - "reasoning_item id={id} encrypted_bytes={encrypted_bytes}" - ))); - } - Ok(ChunkType::Text(text)) => { - emitted_non_replayable_output = true; - last_assistant_message_phase = current_assistant_message_phase; - accumulated_text.push_str(&text); - let _ = tx_loop.send(ChunkType::Text(text)); - } - Ok(ChunkType::Reasoning(reasoning)) => { - emitted_non_replayable_output = true; - accumulated_reasoning.push_str(&reasoning); - let _ = tx_loop.send(ChunkType::Reasoning(reasoning)); - } - Ok(ChunkType::ToolCall(json_str)) => { - emitted_non_replayable_output = true; - has_tool_call = true; - let _ = tx_loop.send(ChunkType::ToolCall(json_str.clone())); - if let Err(err) = tool_call_accumulator.ingest(&json_str) { - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; + Ok(ChunkType::ResponseCompleted { + end_turn, + reasoning_items, + doom_loop_triggers, + }) => { + saw_terminal_event = true; + response_end_turn = end_turn; + for item in reasoning_items { + merge_reasoning_replay_item(&mut reasoning_replay_items, item); + } + if !doom_loop_triggers.is_empty() { + if doom_loop_triggers + .iter() + .any(|trigger| doom_loop_trigger_is_confident(trigger)) + { + server_doom_loop = true; + } + let _ = tx_loop.send(ChunkType::Metadata(format!( + "doom_loop_check triggers={}", + doom_loop_triggers.join(",") + ))); + } + let _ = tx_loop.send(ChunkType::Metadata(format!( + "response.completed end_turn={end_turn:?} reasoning_items={}", + reasoning_replay_items.len() + ))); } - } - Ok(ChunkType::ProviderToolCall(payload)) => { - // Hosted / server-side tools: forward for UI only. - emitted_non_replayable_output = true; - let _ = tx_loop.send(ChunkType::ProviderToolCall(payload)); - } - Ok(ChunkType::End { reason }) => { - // Processed internally — NOT forwarded to tx_loop. - // Forwarding End would cause relay_stream_to_sender - // to return Ended prematurely, dropping the channel - // before tool execution / subsequent steps. - saw_terminal_event = true; - if let Some(reason) = reason { - let label = reason.label().to_string(); - provider_finish_reason = Some(reason); + Ok(ChunkType::ReasoningItem(item)) => { + let id = item.id.clone().unwrap_or_default(); + let encrypted_bytes = item + .encrypted_content + .as_ref() + .map(String::len) + .unwrap_or(0); + merge_reasoning_replay_item(&mut reasoning_replay_items, item); let _ = tx_loop.send(ChunkType::Metadata(format!( - "provider_finish_reason={label}" + "reasoning_item id={id} encrypted_bytes={encrypted_bytes}" ))); } - } - Ok(ChunkType::Metadata(msg)) => { - if doom_loop_metadata_is_confident(&msg) { - server_doom_loop = true; + Ok(ChunkType::Text(text)) => { + emitted_non_replayable_output = true; + last_assistant_message_phase = current_assistant_message_phase; + accumulated_text.push_str(&text); + let _ = tx_loop.send(ChunkType::Text(text)); } - let _ = tx_loop.send(ChunkType::Metadata(msg)); - } - Ok(ChunkType::Warning(msg)) => { - let _ = tx_loop.send(ChunkType::Warning(msg)); - } - Ok(ChunkType::Retry(status)) => { - let _ = tx_loop.send(ChunkType::Retry(status)); - } - Ok(ChunkType::StreamRollback { .. }) => {} - Ok(ChunkType::RetryableFailure(retry_error)) => { - if attempt <= PROVIDER_STEP_MAX_RETRIES { - rollback_provider_attempt( - &tx_loop, - &mut accumulated_text, - &mut accumulated_reasoning, - &mut reasoning_replay_items, - &mut has_tool_call, - &mut tool_call_accumulator, - &mut saw_terminal_event, - &mut response_end_turn, - &mut provider_finish_reason, - &mut last_assistant_message_phase, - &mut current_assistant_message_phase, - &mut emitted_non_replayable_output, - ); - if !emit_retry_and_sleep( - &tx_loop, - step_idx, - attempt, - &retry_error, - cancel_token.as_ref(), - ) - .await - { - let err = "Streaming cancelled by user".to_string(); + Ok(ChunkType::Reasoning(reasoning)) => { + emitted_non_replayable_output = true; + accumulated_reasoning.push_str(&reasoning); + let _ = tx_loop.send(ChunkType::Reasoning(reasoning)); + } + Ok(ChunkType::ToolCall(json_str)) => { + emitted_non_replayable_output = true; + has_tool_call = true; + let _ = tx_loop.send(ChunkType::ToolCall(json_str.clone())); + if let Err(err) = tool_call_accumulator.ingest(&json_str) { let _ = tx_loop.send(ChunkType::Failed(err.clone())); *stop_reason_arc.lock().await = Some(StopReason::Error(err)); return; } - attempt += 1; - stream = match open_provider_stream_with_retries( - &provider_clone, - ¤t_messages, - &tools, - &headers, - &tx_loop, - step_idx, - &mut attempt, - cancel_token.as_ref(), - ) - .await - { - Ok(stream) => stream, - Err(error) => { - let err = provider_step_error_message( - step_idx, - current_messages.len(), - tools.len(), - &step_summary, - error, - ); + } + Ok(ChunkType::ProviderToolCall(payload)) => { + // Hosted / server-side tools: forward for UI only. + emitted_non_replayable_output = true; + had_provider_tool_call = true; + let _ = tx_loop.send(ChunkType::ProviderToolCall(payload)); + } + Ok(ChunkType::End { reason }) => { + // Processed internally — NOT forwarded to tx_loop. + // Forwarding End would cause relay_stream_to_sender + // to return Ended prematurely, dropping the channel + // before tool execution / subsequent steps. + saw_terminal_event = true; + if let Some(reason) = reason { + let label = reason.label().to_string(); + provider_finish_reason = Some(reason); + let _ = tx_loop.send(ChunkType::Metadata(format!( + "provider_finish_reason={label}" + ))); + } + } + Ok(ChunkType::Metadata(msg)) => { + if doom_loop_metadata_is_confident(&msg) { + server_doom_loop = true; + } + let _ = tx_loop.send(ChunkType::Metadata(msg)); + } + Ok(ChunkType::Warning(msg)) => { + let _ = tx_loop.send(ChunkType::Warning(msg)); + } + Ok(ChunkType::Retry(status)) => { + let _ = tx_loop.send(ChunkType::Retry(status)); + } + Ok(ChunkType::StreamRollback { .. }) => {} + Ok(ChunkType::RetryableFailure(retry_error)) => { + if attempt <= PROVIDER_STEP_MAX_RETRIES { + rollback_provider_attempt( + &tx_loop, + &mut accumulated_text, + &mut accumulated_reasoning, + &mut reasoning_replay_items, + &mut has_tool_call, + &mut tool_call_accumulator, + &mut saw_terminal_event, + &mut response_end_turn, + &mut provider_finish_reason, + &mut last_assistant_message_phase, + &mut current_assistant_message_phase, + &mut emitted_non_replayable_output, + &mut had_provider_tool_call, + ); + if !emit_retry_and_sleep( + &tx_loop, + step_idx, + attempt, + &retry_error, + cancel_token.as_ref(), + ) + .await + { + let err = "Streaming cancelled by user".to_string(); let _ = tx_loop.send(ChunkType::Failed(err.clone())); *stop_reason_arc.lock().await = Some(StopReason::Error(err)); return; } - }; - continue; - } + attempt += 1; + stream = match open_provider_stream_with_retries( + &provider_clone, + ¤t_messages, + &tools, + &headers, + &tx_loop, + step_idx, + &mut attempt, + cancel_token.as_ref(), + ) + .await + { + Ok(stream) => stream, + Err(error) => { + let err = provider_step_error_message( + step_idx, + current_messages.len(), + tools.len(), + &step_summary, + error, + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = + Some(StopReason::Error(err)); + return; + } + }; + continue; + } - let err = format!( - "Provider stream failed after {} retries: {}", - PROVIDER_STEP_MAX_RETRIES, retry_error.message - ); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; - } - Ok(ChunkType::Incomplete(msg)) => { - let err = format!("Provider response incomplete: {}", msg); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; - } - Ok(ChunkType::Failed(err)) => { - let retry_error = RetryError::from_message(err.clone()); - if crate::retry::retryable(&retry_error) - && attempt <= PROVIDER_STEP_MAX_RETRIES - { - rollback_provider_attempt( - &tx_loop, - &mut accumulated_text, - &mut accumulated_reasoning, - &mut reasoning_replay_items, - &mut has_tool_call, - &mut tool_call_accumulator, - &mut saw_terminal_event, - &mut response_end_turn, - &mut provider_finish_reason, - &mut last_assistant_message_phase, - &mut current_assistant_message_phase, - &mut emitted_non_replayable_output, + let err = format!( + "Provider stream failed after {} retries: {}", + PROVIDER_STEP_MAX_RETRIES, retry_error.message ); - if !emit_retry_and_sleep( - &tx_loop, - step_idx, - attempt, - &retry_error, - cancel_token.as_ref(), - ) - .await - { - let err = "Streaming cancelled by user".to_string(); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; - } - attempt += 1; - stream = match open_provider_stream_with_retries( - &provider_clone, - ¤t_messages, - &tools, - &headers, - &tx_loop, - step_idx, - &mut attempt, - cancel_token.as_ref(), - ) - .await + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + Ok(ChunkType::Incomplete(msg)) => { + let err = format!("Provider response incomplete: {}", msg); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + Ok(ChunkType::Failed(err)) => { + let retry_error = RetryError::from_message(err.clone()); + if crate::retry::retryable(&retry_error) + && attempt <= PROVIDER_STEP_MAX_RETRIES { - Ok(stream) => stream, - Err(error) => { - let err = provider_step_error_message( - step_idx, - current_messages.len(), - tools.len(), - &step_summary, - error, - ); + rollback_provider_attempt( + &tx_loop, + &mut accumulated_text, + &mut accumulated_reasoning, + &mut reasoning_replay_items, + &mut has_tool_call, + &mut tool_call_accumulator, + &mut saw_terminal_event, + &mut response_end_turn, + &mut provider_finish_reason, + &mut last_assistant_message_phase, + &mut current_assistant_message_phase, + &mut emitted_non_replayable_output, + &mut had_provider_tool_call, + ); + if !emit_retry_and_sleep( + &tx_loop, + step_idx, + attempt, + &retry_error, + cancel_token.as_ref(), + ) + .await + { + let err = "Streaming cancelled by user".to_string(); let _ = tx_loop.send(ChunkType::Failed(err.clone())); *stop_reason_arc.lock().await = Some(StopReason::Error(err)); return; } - }; - continue; - } + attempt += 1; + stream = match open_provider_stream_with_retries( + &provider_clone, + ¤t_messages, + &tools, + &headers, + &tx_loop, + step_idx, + &mut attempt, + cancel_token.as_ref(), + ) + .await + { + Ok(stream) => stream, + Err(error) => { + let err = provider_step_error_message( + step_idx, + current_messages.len(), + tools.len(), + &step_summary, + error, + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = + Some(StopReason::Error(err)); + return; + } + }; + continue; + } - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; - } - Ok(ChunkType::Start) => { - let _ = tx_loop.send(ChunkType::Start); - } - Ok(ChunkType::NotSupported(msg)) => { - let _ = tx_loop.send(ChunkType::NotSupported(msg)); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + Ok(ChunkType::Start) => { + let _ = tx_loop.send(ChunkType::Start); + } + Ok(ChunkType::NotSupported(msg)) => { + let _ = tx_loop.send(ChunkType::NotSupported(msg)); + } + Err(e) => match retry_error_from_provider_error(e) { + Ok(retry_error) if attempt <= PROVIDER_STEP_MAX_RETRIES => { + rollback_provider_attempt( + &tx_loop, + &mut accumulated_text, + &mut accumulated_reasoning, + &mut reasoning_replay_items, + &mut has_tool_call, + &mut tool_call_accumulator, + &mut saw_terminal_event, + &mut response_end_turn, + &mut provider_finish_reason, + &mut last_assistant_message_phase, + &mut current_assistant_message_phase, + &mut emitted_non_replayable_output, + &mut had_provider_tool_call, + ); + if !emit_retry_and_sleep( + &tx_loop, + step_idx, + attempt, + &retry_error, + cancel_token.as_ref(), + ) + .await + { + let err = "Streaming cancelled by user".to_string(); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + attempt += 1; + stream = match open_provider_stream_with_retries( + &provider_clone, + ¤t_messages, + &tools, + &headers, + &tx_loop, + step_idx, + &mut attempt, + cancel_token.as_ref(), + ) + .await + { + Ok(stream) => stream, + Err(error) => { + let err = provider_step_error_message( + step_idx, + current_messages.len(), + tools.len(), + &step_summary, + error, + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = + Some(StopReason::Error(err)); + return; + } + }; + continue; + } + Ok(retry_error) => { + let err = format!( + "Provider stream failed after {} retries: {}", + PROVIDER_STEP_MAX_RETRIES, retry_error.message + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + Err(error) => { + let err = error.to_string(); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + }, } - Err(e) => match retry_error_from_provider_error(e) { - Ok(retry_error) if attempt <= PROVIDER_STEP_MAX_RETRIES => { + } + + if !saw_terminal_event { + let err = + "Provider stream ended without a terminal completion event".to_string(); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + + // Grok Build sampler: a completed response with no visible + // content and no tool calls is Empty (reasoning-only or + // no_visible_content). Retry the same request; do not accept + // it as a finished turn or append it to the conversation. + // `.devrefs/.../xai-grok-sampler/src/actor/request_task.rs` + // `AttemptOutcome::Empty` + `ConversationResponse::empty_reason`. + // Content-filter empties are deterministic and must not retry. + if !has_tool_call + && !had_provider_tool_call + && accumulated_text.trim().is_empty() + && !matches!(provider_finish_reason, Some(FinishReason::ContentFilter)) + { + let had_reasoning = !accumulated_reasoning.is_empty() + || reasoning_replay_items.iter().any(|item| !item.is_empty()); + let empty_reason = if had_reasoning { + "reasoning_only" + } else { + "no_visible_content" + }; + let _ = tx_loop.send(ChunkType::Metadata(format!( + "empty_response reason={empty_reason} had_reasoning={had_reasoning} content_len=0 tool_call_count=0 attempt={attempt}" + ))); + // Grok Build: doom outranks empty. A reasoning-only + // sample with tail_repetition@thinking is resampled + // with the recovery reminder, not the same request. + if server_doom_loop { + if let Some(reminder_text) = doom_loop.begin_recovery() { + let _ = tx_loop.send(ChunkType::Metadata(format!( + "doom_loop_detected step={} recoveries={} empty_reason={empty_reason}", + step_idx, doom_loop.recoveries, + ))); rollback_provider_attempt( &tx_loop, &mut accumulated_text, @@ -494,21 +632,15 @@ pub async fn stream_with_tools( &mut last_assistant_message_phase, &mut current_assistant_message_phase, &mut emitted_non_replayable_output, + &mut had_provider_tool_call, ); - if !emit_retry_and_sleep( - &tx_loop, - step_idx, - attempt, - &retry_error, - cancel_token.as_ref(), - ) - .await + if current_messages + .last() + .is_some_and(is_injected_system_reminder) { - let err = "Streaming cancelled by user".to_string(); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; + current_messages.pop(); } + current_messages.push(Message::user(reminder_text)); attempt += 1; stream = match open_provider_stream_with_retries( &provider_clone, @@ -536,32 +668,80 @@ pub async fn stream_with_tools( return; } }; - continue; - } - Ok(retry_error) => { - let err = format!( - "Provider stream failed after {} retries: {}", - PROVIDER_STEP_MAX_RETRIES, retry_error.message - ); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; + continue 'resample; } - Err(error) => { - let err = error.to_string(); + } + if attempt <= PROVIDER_STEP_MAX_RETRIES { + let retry_error = + RetryError::from_message(format!("empty response: {empty_reason}")); + rollback_provider_attempt( + &tx_loop, + &mut accumulated_text, + &mut accumulated_reasoning, + &mut reasoning_replay_items, + &mut has_tool_call, + &mut tool_call_accumulator, + &mut saw_terminal_event, + &mut response_end_turn, + &mut provider_finish_reason, + &mut last_assistant_message_phase, + &mut current_assistant_message_phase, + &mut emitted_non_replayable_output, + &mut had_provider_tool_call, + ); + if !emit_retry_and_sleep( + &tx_loop, + step_idx, + attempt, + &retry_error, + cancel_token.as_ref(), + ) + .await + { + let err = "Streaming cancelled by user".to_string(); let _ = tx_loop.send(ChunkType::Failed(err.clone())); *stop_reason_arc.lock().await = Some(StopReason::Error(err)); return; } - }, + attempt += 1; + stream = match open_provider_stream_with_retries( + &provider_clone, + ¤t_messages, + &tools, + &headers, + &tx_loop, + step_idx, + &mut attempt, + cancel_token.as_ref(), + ) + .await + { + Ok(stream) => stream, + Err(error) => { + let err = provider_step_error_message( + step_idx, + current_messages.len(), + tools.len(), + &step_summary, + error, + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; + } + }; + continue 'resample; + } + let err = format!( + "Provider stream failed after {} retries: empty response: {empty_reason}", + PROVIDER_STEP_MAX_RETRIES + ); + let _ = tx_loop.send(ChunkType::Failed(err.clone())); + *stop_reason_arc.lock().await = Some(StopReason::Error(err)); + return; } - } - if !saw_terminal_event { - let err = "Provider stream ended without a terminal completion event".to_string(); - let _ = tx_loop.send(ChunkType::Failed(err.clone())); - *stop_reason_arc.lock().await = Some(StopReason::Error(err)); - return; + break 'resample; } // Responses order: reasoning siblings, then assistant text (if any), @@ -840,6 +1020,7 @@ fn rollback_provider_attempt( last_assistant_message_phase: &mut Option, current_assistant_message_phase: &mut Option, emitted_non_replayable_output: &mut bool, + had_provider_tool_call: &mut bool, ) { if !accumulated_text.is_empty() || !accumulated_reasoning.is_empty() { let _ = tx.send(ChunkType::StreamRollback { @@ -859,6 +1040,7 @@ fn rollback_provider_attempt( *last_assistant_message_phase = None; *current_assistant_message_phase = None; *emitted_non_replayable_output = false; + *had_provider_tool_call = false; } /// Prune older tool outputs in the live multi-step transcript before each @@ -2120,30 +2302,61 @@ mod tests { } #[derive(Debug, Clone)] - struct RecoveringToolFailureProvider { + struct ReasoningOnlyCompletedProvider { requests: Arc, - observed_follow_up: Arc>>, } #[derive(Debug, Clone)] - struct RecoveringRateLimitProvider { + struct ReasoningOnlyDoomLoopProvider { requests: Arc, + saw_reminder: Arc, } #[derive(Debug, Clone)] - struct RecoveringPartialStreamProvider { + struct NoVisibleContentCompletedProvider { requests: Arc, } - #[async_trait] - impl Provider for RecoveringRateLimitProvider { - fn name(&self) -> &str { - "test" - } - - fn model_name(&self) -> &str { - "test" - } + #[derive(Debug, Clone)] + struct AlwaysEmptyCompletedProvider { + requests: Arc, + } + + #[derive(Debug, Clone)] + struct ContentFilterEmptyProvider { + requests: Arc, + } + + #[derive(Debug, Clone)] + struct HostedToolOnlyCompletedProvider { + requests: Arc, + } + + #[derive(Debug, Clone)] + struct RecoveringToolFailureProvider { + requests: Arc, + observed_follow_up: Arc>>, + } + + #[derive(Debug, Clone)] + struct RecoveringRateLimitProvider { + requests: Arc, + } + + #[derive(Debug, Clone)] + struct RecoveringPartialStreamProvider { + requests: Arc, + } + + #[async_trait] + impl Provider for RecoveringRateLimitProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } async fn stream_text( &self, @@ -2649,6 +2862,209 @@ mod tests { } } + #[async_trait] + impl Provider for ReasoningOnlyCompletedProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + _messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + let request = self.requests.fetch_add(1, Ordering::SeqCst); + let chunks = match request { + 0 => vec![ + Ok(ChunkType::Reasoning( + "Let me view the screenshots.".to_string(), + )), + Ok(ChunkType::response_completed(None)), + ], + 1 => vec![ + Ok(ChunkType::ToolCall( + r#"[{"index":0,"id":"call_list","type":"function","function":{"name":"list","arguments":"{\"path\":\".\"}"}}]"# + .to_string(), + )), + Ok(ChunkType::End { + reason: Some(FinishReason::ToolCalls), + }), + ], + _ => vec![ + Ok(ChunkType::Text("Done.".to_string())), + Ok(ChunkType::End { + reason: Some(FinishReason::Stop), + }), + ], + }; + + Ok(Box::pin(futures::stream::iter(chunks))) + } + } + + #[async_trait] + impl Provider for ReasoningOnlyDoomLoopProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + let request = self.requests.fetch_add(1, Ordering::SeqCst); + if messages.iter().any(|message| { + matches!( + message, + Message::User(user) if user.content.starts_with("") + ) + }) { + self.saw_reminder.store(true, Ordering::SeqCst); + } + let chunks = match request { + 0 => vec![ + Ok(ChunkType::Reasoning("looping thought".to_string())), + Ok(ChunkType::ResponseCompleted { + end_turn: None, + reasoning_items: Vec::new(), + doom_loop_triggers: vec!["tail_repetition:8@thinking".to_string()], + }), + ], + 1 => vec![ + Ok(ChunkType::ToolCall( + r#"[{"index":0,"id":"call_list","type":"function","function":{"name":"list","arguments":"{\"path\":\".\"}"}}]"# + .to_string(), + )), + Ok(ChunkType::End { + reason: Some(FinishReason::ToolCalls), + }), + ], + _ => vec![ + Ok(ChunkType::Text("Done.".to_string())), + Ok(ChunkType::End { + reason: Some(FinishReason::Stop), + }), + ], + }; + + Ok(Box::pin(futures::stream::iter(chunks))) + } + } + + #[async_trait] + impl Provider for NoVisibleContentCompletedProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + _messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + let request = self.requests.fetch_add(1, Ordering::SeqCst); + let chunks = if request == 0 { + vec![Ok(ChunkType::response_completed(None))] + } else { + vec![ + Ok(ChunkType::Text("Done.".to_string())), + Ok(ChunkType::End { + reason: Some(FinishReason::Stop), + }), + ] + }; + Ok(Box::pin(futures::stream::iter(chunks))) + } + } + + #[async_trait] + impl Provider for AlwaysEmptyCompletedProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + _messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + self.requests.fetch_add(1, Ordering::SeqCst); + Ok(Box::pin(futures::stream::iter(vec![Ok( + ChunkType::response_completed(None), + )]))) + } + } + + #[async_trait] + impl Provider for ContentFilterEmptyProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + _messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + self.requests.fetch_add(1, Ordering::SeqCst); + Ok(Box::pin(futures::stream::iter(vec![Ok(ChunkType::End { + reason: Some(FinishReason::ContentFilter), + })]))) + } + } + + #[async_trait] + impl Provider for HostedToolOnlyCompletedProvider { + fn name(&self) -> &str { + "test" + } + + fn model_name(&self) -> &str { + "test" + } + + async fn stream_text( + &self, + _messages: &[Message], + _tools: &[Tool], + _headers: &HashMap, + ) -> crate::error::Result { + self.requests.fetch_add(1, Ordering::SeqCst); + Ok(Box::pin(futures::stream::iter(vec![ + Ok(ChunkType::ProviderToolCall( + r#"{"id":"hs_1","name":"web_search","status":"completed"}"#.to_string(), + )), + Ok(ChunkType::response_completed(None)), + ]))) + } + } + #[async_trait] impl Provider for RecoveringToolFailureProvider { fn name(&self) -> &str { @@ -3479,6 +3895,280 @@ mod tests { assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); } + fn list_tool(executions: Arc) -> Tool { + Tool::builder() + .name("list") + .description("list files") + .input_schema(Schema::from(true)) + .execute(ToolExecute::new(move |_input| { + let executions = executions.clone(); + async move { + executions.fetch_add(1, Ordering::SeqCst); + Ok("package.json".to_string()) + } + })) + .build() + .unwrap() + } + + #[tokio::test(start_paused = true)] + async fn resamples_reasoning_only_completed_response_like_grok_build() { + let provider = ReasoningOnlyCompletedProvider { + requests: Arc::new(AtomicUsize::new(0)), + }; + let executions = Arc::new(AtomicUsize::new(0)); + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("continue")], + vec![list_tool(executions.clone())], + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut text = String::new(); + let mut visible_reasoning = String::new(); + let mut empty_logged = false; + let mut saw_rollback = false; + let mut saw_retry = false; + while let Some(chunk) = response.stream.next().await { + match chunk { + ChunkType::Text(delta) => text.push_str(&delta), + ChunkType::Reasoning(reasoning) => visible_reasoning.push_str(&reasoning), + ChunkType::StreamRollback { reasoning, .. } => { + saw_rollback = true; + if visible_reasoning.ends_with(&reasoning) { + visible_reasoning.truncate(visible_reasoning.len() - reasoning.len()); + } + } + ChunkType::Retry(_) => saw_retry = true, + ChunkType::Metadata(message) + if message.contains("empty_response reason=reasoning_only") => + { + empty_logged = true; + } + _ => {} + } + } + + assert_eq!(text, "Done."); + assert!(empty_logged); + assert!(saw_rollback); + assert!(saw_retry); + assert_eq!(provider.requests.load(Ordering::SeqCst), 3); + assert_eq!(executions.load(Ordering::SeqCst), 1); + assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); + let durable = response.messages().await; + assert!( + durable.iter().all(|message| match message { + Message::User(user) => !user.content.starts_with(""), + _ => true, + }), + "empty resample must not inject a user reminder: {durable:?}" + ); + } + + #[tokio::test(start_paused = true)] + async fn resamples_no_visible_content_completed_response() { + let provider = NoVisibleContentCompletedProvider { + requests: Arc::new(AtomicUsize::new(0)), + }; + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("continue")], + vec![list_tool(Arc::new(AtomicUsize::new(0)))], + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut text = String::new(); + let mut empty_logged = false; + while let Some(chunk) = response.stream.next().await { + match chunk { + ChunkType::Text(delta) => text.push_str(&delta), + ChunkType::Metadata(message) + if message.contains("empty_response reason=no_visible_content") => + { + empty_logged = true; + } + _ => {} + } + } + + assert_eq!(text, "Done."); + assert!(empty_logged); + assert_eq!(provider.requests.load(Ordering::SeqCst), 2); + assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); + } + + #[tokio::test(start_paused = true)] + async fn empty_response_exhaustion_fails_instead_of_finish() { + let provider = AlwaysEmptyCompletedProvider { + requests: Arc::new(AtomicUsize::new(0)), + }; + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("continue")], + Vec::new(), + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut failed = Vec::new(); + while let Some(chunk) = response.stream.next().await { + if let ChunkType::Failed(message) = chunk { + failed.push(message); + } + } + + assert!( + failed.iter().any(|message| { + message.contains("empty response: no_visible_content") + && message.contains("after 10 retries") + }), + "expected empty-response exhaustion, got {failed:?}" + ); + assert_eq!(provider.requests.load(Ordering::SeqCst), 11); + assert!(matches!( + response.stop_reason().await, + Some(StopReason::Error(message)) + if message.contains("empty response: no_visible_content") + )); + } + + #[tokio::test] + async fn content_filter_empty_does_not_resample() { + let provider = ContentFilterEmptyProvider { + requests: Arc::new(AtomicUsize::new(0)), + }; + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("continue")], + Vec::new(), + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut retries = 0usize; + let mut empty_logged = false; + while let Some(chunk) = response.stream.next().await { + match chunk { + ChunkType::Retry(_) => retries += 1, + ChunkType::Metadata(message) if message.contains("empty_response") => { + empty_logged = true; + } + _ => {} + } + } + + assert!(!empty_logged); + assert_eq!(retries, 0); + assert_eq!(provider.requests.load(Ordering::SeqCst), 1); + assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); + } + + #[tokio::test] + async fn hosted_tool_only_completed_response_is_not_empty() { + let provider = HostedToolOnlyCompletedProvider { + requests: Arc::new(AtomicUsize::new(0)), + }; + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("search the web")], + Vec::new(), + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut hosted = 0usize; + let mut retries = 0usize; + let mut empty_logged = false; + while let Some(chunk) = response.stream.next().await { + match chunk { + ChunkType::ProviderToolCall(_) => hosted += 1, + ChunkType::Retry(_) => retries += 1, + ChunkType::Metadata(message) if message.contains("empty_response") => { + empty_logged = true; + } + _ => {} + } + } + + assert_eq!(hosted, 1); + assert!(!empty_logged); + assert_eq!(retries, 0); + assert_eq!(provider.requests.load(Ordering::SeqCst), 1); + assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); + } + + #[tokio::test] + async fn reasoning_only_with_thinking_doom_loop_resamples_with_reminder() { + let saw_reminder = Arc::new(AtomicBool::new(false)); + let provider = ReasoningOnlyDoomLoopProvider { + requests: Arc::new(AtomicUsize::new(0)), + saw_reminder: saw_reminder.clone(), + }; + let executions = Arc::new(AtomicUsize::new(0)); + + let mut response = stream_with_tools( + provider.clone(), + vec![Message::user("continue")], + vec![list_tool(executions.clone())], + Some(5), + None, + HashMap::new(), + None, + ) + .await + .unwrap(); + + let mut doom_logged = false; + while let Some(chunk) = response.stream.next().await { + if let ChunkType::Metadata(message) = chunk { + if message.contains("doom_loop_detected") && message.contains("empty_reason=") { + doom_logged = true; + } + } + } + + assert!(doom_logged); + assert!(saw_reminder.load(Ordering::SeqCst)); + assert!(executions.load(Ordering::SeqCst) >= 1); + assert_eq!(response.stop_reason().await, Some(StopReason::Finish)); + let durable = response.messages().await; + assert!( + durable.iter().all(|message| match message { + Message::User(user) => !user.content.starts_with(""), + _ => true, + }), + "doom reminder is request-only: {durable:?}" + ); + } + #[tokio::test] async fn continues_once_after_phase_less_end_turn_without_final_phase() { let provider = PhaselessAmbiguousProvider { diff --git a/src/aisdk/retry.rs b/src/aisdk/retry.rs index 14db94b2..fb23c351 100644 --- a/src/aisdk/retry.rs +++ b/src/aisdk/retry.rs @@ -126,6 +126,9 @@ pub fn retryable(error: &RetryError) -> bool { || lower.contains("error sending request") || lower.contains("server_error") || lower.contains("too_many_requests") + || lower.contains("empty response") + || lower.contains("reasoning_only") + || lower.contains("no_visible_content") } pub fn retry_message(error: &RetryError) -> String { diff --git a/src/llm/xai_build.rs b/src/llm/xai_build.rs index c99c2b0f..629595a3 100644 --- a/src/llm/xai_build.rs +++ b/src/llm/xai_build.rs @@ -17,6 +17,10 @@ const PROTOCOL_VERSION_FALLBACK: &str = "0.2.111"; /// (`DoomLoopRecoveryPolicy::DEFAULT_RECOVERY_WINDOW_TOKENS`). const DOOM_LOOP_CHECK_HEADER: &str = "x-grok-doom-loop-check"; const DOOM_LOOP_CHECK_WINDOW_TOKENS: &str = "1024"; +/// Grok Build sibling header; default is the 64-token primitive minimum. +/// `.devrefs/references/xai-org/grok-build/crates/codegen/xai-grok-sampling-types/src/doom_loop.rs` +const EXACT_REPETITION_CHECK_HEADER: &str = "x-grok-exact-repetition-check"; +const EXACT_REPETITION_MIN_TOKENS: &str = "64"; const VERSION_CACHE_TTL: Duration = Duration::from_secs(24 * 60 * 60); const VERSION_RETRY_TTL: Duration = Duration::from_secs(5 * 60); const VERSION_URLS: &[&str] = &[ @@ -115,6 +119,10 @@ fn request_overrides_with_version( DOOM_LOOP_CHECK_HEADER.to_string(), DOOM_LOOP_CHECK_WINDOW_TOKENS.to_string(), ); + headers.insert( + EXACT_REPETITION_CHECK_HEADER.to_string(), + EXACT_REPETITION_MIN_TOKENS.to_string(), + ); RequestOverrides { api_key: oauth_access, @@ -475,6 +483,13 @@ mod tests { .map(String::as_str), Some(super::DOOM_LOOP_CHECK_WINDOW_TOKENS) ); + assert_eq!( + overrides + .headers + .get(super::EXACT_REPETITION_CHECK_HEADER) + .map(String::as_str), + Some(super::EXACT_REPETITION_MIN_TOKENS) + ); } #[test]