| 1 | //! Streaming response state and guardrails. |
| 2 | //! |
| 3 | //! This module owns the local state used while decoding one model stream: |
| 4 | //! content block kind tracking, streamed tool-use buffers, transparent retry |
| 5 | //! policy, and scrubbers for text that looks like a forged tool-call wrapper. |
| 6 | |
| 7 | use crate::core::events::TurnOutcomeStatus; |
| 8 | use codewhale_models::ToolCaller; |
| 9 | use std::time::Duration; |
| 10 | |
| 11 | /// A send that did not enter the existing event queue. Cancellation is not |
| 12 | /// evidence that the consumer closed, and neither is user-visible delivery. |
| 13 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 14 | pub(super) enum EventSendError { |
| 15 | Cancelled, |
| 16 | Closed, |
| 17 | } |
| 18 | |
| 19 | #[derive(Debug, Clone, Copy, PartialEq, Eq)] |
| 20 | pub(super) enum EventReservationPolicy { |
| 21 | /// Cancellation wins before admitting work or a new stream observation. |
| 22 | Strict, |
| 23 | /// Preserve a completed observation when capacity is already available. |
| 24 | Receipt, |
| 25 | } |
| 26 | |
| 27 | /// Reserve from the one event queue. Lifecycle and terminal handoff permits |
| 28 | /// are held locally; ordinary receipts consume theirs immediately. Receipts |
| 29 | /// may enter available capacity after cancellation, but cancellation always |
| 30 | /// releases a wait on a full queue. Idle sends have no turn token to cancel. |
| 31 | pub(super) async fn reserve_event_capacity( |
| 32 | tx: &tokio::sync::mpsc::Sender<super::Event>, |
| 33 | cancel: Option<&tokio_util::sync::CancellationToken>, |
| 34 | policy: EventReservationPolicy, |
| 35 | ) -> Result<tokio::sync::mpsc::OwnedPermit<super::Event>, EventSendError> { |
| 36 | if policy == EventReservationPolicy::Receipt { |
| 37 | match tx.clone().try_reserve_owned() { |
| 38 | Ok(permit) => return Ok(permit), |
| 39 | Err(tokio::sync::mpsc::error::TrySendError::Closed(_)) => { |
| 40 | return Err(EventSendError::Closed); |
| 41 | } |
| 42 | Err(tokio::sync::mpsc::error::TrySendError::Full(_)) => {} |
| 43 | } |
| 44 | } |
| 45 | let reserve = tx.clone().reserve_owned(); |
| 46 | match cancel { |
| 47 | Some(cancel) => tokio::select! { |
| 48 | biased; |
| 49 | () = cancel.cancelled() => Err(EventSendError::Cancelled), |
| 50 | result = reserve => result.map_err(|_| EventSendError::Closed), |
| 51 | }, |
| 52 | None => reserve.await.map_err(|_| EventSendError::Closed), |
| 53 | } |
| 54 | } |
| 55 | |
| 56 | /// Selective quiet affects only attempt observations; completed summaries and |
| 57 | /// counters survive. Admission and delivery use the same existing queue guard. |
| 58 | async fn emit_retry_status( |
| 59 | tx: &tokio::sync::mpsc::Sender<super::Event>, |
| 60 | cancel: &tokio_util::sync::CancellationToken, |
| 61 | quiet: bool, |
| 62 | message: String, |
| 63 | ) -> Result<(), EventSendError> { |
| 64 | let attempt = message.starts_with("Retry attempt:"); |
| 65 | if quiet && attempt { |
| 66 | return Ok(()); |
| 67 | } |
| 68 | let policy = if attempt { |
| 69 | EventReservationPolicy::Strict |
| 70 | } else { |
| 71 | EventReservationPolicy::Receipt |
| 72 | }; |
| 73 | let permit = reserve_event_capacity(tx, Some(cancel), policy).await?; |
| 74 | permit.send(super::Event::status(message)); |
| 75 | Ok(()) |
| 76 | } |
| 77 | |
| 78 | impl super::Engine { |
| 79 | pub(super) fn request_retry_observation(&self) -> crate::llm_client::RequestRetryObservation { |
| 80 | let tx = self.tx_event.clone(); |
| 81 | let cancel = self.cancel_token.clone(); |
| 82 | let quiet = self.api_config.notifications_config().quiet; |
| 83 | crate::llm_client::RequestRetryObservation { |
| 84 | retries: std::sync::Arc::new(std::sync::atomic::AtomicU32::new(0)), |
| 85 | emit: std::sync::Arc::new(move |message| { |
| 86 | let tx = tx.clone(); |
| 87 | let cancel = cancel.clone(); |
| 88 | Box::pin(async move { |
| 89 | let _ = emit_retry_status(&tx, &cancel, quiet, message).await; |
| 90 | }) |
| 91 | }), |
| 92 | } |
| 93 | } |
| 94 | |
| 95 | // The existing turn owns these cumulative counters. Closing observations |
| 96 | // describe its actual outcome, without attributing a later tool failure |
| 97 | // to an earlier provider response or changing any recovery budget. |
| 98 | pub(super) async fn send_answer_retry_summary( |
| 99 | &self, |
| 100 | diagnostics: &crate::tool_inspection::TurnStopDiagnostics, |
| 101 | status: TurnOutcomeStatus, |
| 102 | ) { |
| 103 | let (prefix, outcome) = match status { |
| 104 | TurnOutcomeStatus::Completed => ("Retry recovery", "turn completed"), |
| 105 | TurnOutcomeStatus::Failed => ("Retry stopped", "turn failed"), |
| 106 | TurnOutcomeStatus::Interrupted => ("Retry interrupted", "turn interrupted"), |
| 107 | }; |
| 108 | for (kind, retries, limit) in [ |
| 109 | ( |
| 110 | "reasoning-only", |
| 111 | diagnostics.reasoning_only_reprompts, |
| 112 | self.config.reasoning_only_max_reprompts, |
| 113 | ), |
| 114 | ( |
| 115 | "empty-stop", |
| 116 | diagnostics.empty_stop_retries, |
| 117 | super::turn_loop::EMPTY_STOP_MAX_RETRIES, |
| 118 | ), |
| 119 | ] { |
| 120 | if retries > 0 { |
| 121 | let _ = self |
| 122 | .send_retry_status(format!( |
| 123 | "{prefix}: {kind} used {retries}/{limit} retries; {outcome}" |
| 124 | )) |
| 125 | .await; |
| 126 | } |
| 127 | } |
| 128 | } |
| 129 | |
| 130 | pub(super) async fn send_retry_status(&self, message: String) -> Result<(), EventSendError> { |
| 131 | emit_retry_status( |
| 132 | &self.tx_event, |
| 133 | &self.cancel_token, |
| 134 | self.api_config.notifications_config().quiet, |
| 135 | message, |
| 136 | ) |
| 137 | .await |
| 138 | } |
| 139 | |
| 140 | /// Stream observations always belong to the decoder's current turn, |
| 141 | /// including direct test/embedding calls that do not enqueue an Op. |
| 142 | pub(super) async fn send_stream_event(&self, event: super::Event) -> bool { |
| 143 | match reserve_event_capacity( |
| 144 | &self.tx_event, |
| 145 | Some(&self.cancel_token), |
| 146 | EventReservationPolicy::Strict, |
| 147 | ) |
| 148 | .await |
| 149 | { |
| 150 | Ok(permit) => { |
| 151 | permit.send(event); |
| 152 | true |
| 153 | } |
| 154 | Err(_) => false, |
| 155 | } |
| 156 | } |
| 157 | |
| 158 | /// Every instance emitter uses the same queue authority. Available |
| 159 | /// capacity preserves post-cancel usage/status receipts; cancellation |
| 160 | /// releases a wait for capacity. An idle refresh must not inherit the |
| 161 | /// token of an earlier interrupted turn. Stream/admission/terminal handoff |
| 162 | /// reservations retain their strict cancellation floor. |
| 163 | pub(super) async fn send_event(&self, event: super::Event) -> Result<(), EventSendError> { |
| 164 | let cancel = self |
| 165 | .turn_controls |
| 166 | .lock() |
| 167 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 168 | .active |
| 169 | .as_ref() |
| 170 | .map(|control| control.cancel.clone()); |
| 171 | let permit = reserve_event_capacity( |
| 172 | &self.tx_event, |
| 173 | cancel.as_ref(), |
| 174 | EventReservationPolicy::Receipt, |
| 175 | ) |
| 176 | .await?; |
| 177 | permit.send(event); |
| 178 | Ok(()) |
| 179 | } |
| 180 | } |
| 181 | |
| 182 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 183 | pub(super) enum ContentBlockKind { |
| 184 | Text, |
| 185 | Thinking, |
| 186 | ToolUse, |
| 187 | } |
| 188 | |
| 189 | #[derive(Debug, Clone)] |
| 190 | pub(super) struct ToolUseState { |
| 191 | pub(super) id: String, |
| 192 | pub(super) execution_id: String, |
| 193 | pub(super) name: String, |
| 194 | pub(super) input: serde_json::Value, |
| 195 | pub(super) caller: Option<ToolCaller>, |
| 196 | /// Google thought signature captured on the tool call; replayed with the |
| 197 | /// assistant tool-call message on later turns. |
| 198 | pub(super) thought_signature: Option<String>, |
| 199 | pub(super) input_buffer: String, |
| 200 | pub(super) input_parse_error: Option<String>, |
| 201 | } |
| 202 | |
| 203 | impl ToolUseState { |
| 204 | pub(super) fn model_call(&self) -> crate::core::events::ModelToolCall { |
| 205 | crate::core::events::ModelToolCall { |
| 206 | provider_id: self.id.clone(), |
| 207 | caller: self.caller.clone(), |
| 208 | thought_signature: self.thought_signature.clone(), |
| 209 | } |
| 210 | } |
| 211 | } |
| 212 | |
| 213 | /// Maximum total bytes of text, reasoning and tool-argument content before aborting the stream. |
| 214 | pub(super) const STREAM_MAX_CONTENT_BYTES: usize = 10 * 1024 * 1024; // 10 MB |
| 215 | /// A response can contain many empty tool starts without spending the byte |
| 216 | /// budget. Bound that batch before any call is retained or admitted. A lower |
| 217 | /// configured per-turn tool budget remains authoritative at execution. |
| 218 | pub(super) const MAX_TOOL_CALLS_PER_RESPONSE: usize = 256; |
| 219 | |
| 220 | pub(super) fn tool_call_limit_error() -> crate::error_taxonomy::ErrorEnvelope { |
| 221 | crate::error_taxonomy::ErrorEnvelope::new( |
| 222 | crate::error_taxonomy::ErrorCategory::InvalidInput, |
| 223 | crate::error_taxonomy::ErrorSeverity::Error, |
| 224 | false, |
| 225 | "response_tool_call_limit", |
| 226 | format!( |
| 227 | "Model response exceeded the maximum of {MAX_TOOL_CALLS_PER_RESPONSE} tool calls; no call from this response was executed" |
| 228 | ), |
| 229 | ) |
| 230 | } |
| 231 | |
| 232 | /// Sanity backstop for total stream wall-clock duration. **Not** a routine |
| 233 | /// kill switch — the stream chunk idle timeout is the primary stall |
| 234 | /// detector. The wall-clock cap is here only to bound pathological cases |
| 235 | /// (e.g. a server that keeps sending heartbeats forever without progress). |
| 236 | /// |
| 237 | /// History: this used to be 300s (5 min) which was too aggressive — V4 |
| 238 | /// thinking turns on hard prompts legitimately exceed 5 minutes wall-clock |
| 239 | /// while still emitting reasoning_content chunks the whole way. Bumped to |
| 240 | /// 30 min in v0.6.6 after long-reasoning turns hit the old cap. Codex defaults to a |
| 241 | /// per-chunk idle of 300s with no wall-clock cap; we keep both layers but |
| 242 | /// give the wall-clock a generous window so it never fires in practice. |
| 243 | pub(super) const STREAM_MAX_DURATION_SECS: u64 = 1800; // 30 minutes (was 300s; #103/#1) |
| 244 | /// Hard cap on consecutive recoverable stream errors before we surface a turn |
| 245 | /// failure. Bumped 3 → 5 in v0.6.7 along with the HTTP/2 keepalive defaults |
| 246 | /// (#103) — keepalive should make spurious decode errors rarer, so we can |
| 247 | /// tolerate a longer streak before giving up on the turn. This is the |
| 248 | /// default; `[tui].stream_max_errors` overrides it (#6700). |
| 249 | pub(super) const MAX_STREAM_ERRORS_BEFORE_FAIL: u32 = 5; |
| 250 | /// Cap on transparent stream-level retries — these only happen when the wire |
| 251 | /// dies before any content was streamed. The user has seen nothing, but |
| 252 | /// provider usage or billing may already exist. Two attempts can ride out a |
| 253 | /// flaky edge node without amplifying real outages (#103). This is the |
| 254 | /// default; `[tui].stream_max_transparent_retries` overrides it (#6700). |
| 255 | pub(super) const MAX_TRANSPARENT_STREAM_RETRIES: u32 = 2; |
| 256 | |
| 257 | /// Decide whether a stream error is eligible for a transparent retry. |
| 258 | /// |
| 259 | /// True only when ALL three conditions hold: |
| 260 | /// 1. No content has been received on the current attempt. Reissuing after |
| 261 | /// visible partial deltas needs a separate recovery policy. This content |
| 262 | /// check is not evidence that the provider consumed or billed zero tokens. |
| 263 | /// 2. We still have transparent-retry budget remaining (`max_attempts`, |
| 264 | /// `[tui].stream_max_transparent_retries`, default |
| 265 | /// [`MAX_TRANSPARENT_STREAM_RETRIES`]). |
| 266 | /// 3. The turn has not been cancelled. |
| 267 | /// |
| 268 | /// Extracted as a pure function so the four #103 retry cases can be exercised |
| 269 | /// in unit tests without booting the full engine state machine. |
| 270 | pub(super) fn should_transparently_retry_stream( |
| 271 | any_content_received: bool, |
| 272 | transparent_attempts: u32, |
| 273 | max_attempts: u32, |
| 274 | cancelled: bool, |
| 275 | ) -> bool { |
| 276 | !any_content_received && transparent_attempts < max_attempts && !cancelled |
| 277 | } |
| 278 | |
| 279 | /// Default budget for re-issuing the whole request after a dead stream. |
| 280 | /// Shared by the nothing-streamed outer retry (#103 Phase 3), the |
| 281 | /// sleep-resume retry (#2990), the network-drop resumes, and stream-open |
| 282 | /// failures (#6699). Overridable via `[tui].stream_max_resumes` (#6700). |
| 283 | pub(super) const MAX_STREAM_RETRIES: u32 = 3; |
| 284 | |
| 285 | /// Typed, engine-internal state for one mid-stream drop recovery. |
| 286 | /// |
| 287 | /// This enum **is** the retry mechanism. A resumed turn used to append a |
| 288 | /// synthetic `[runtime]` *user* message to the persisted conversation, which |
| 289 | /// polluted the transcript and — when only hidden reasoning had streamed — |
| 290 | /// promised a preserved partial answer that never existed (0.9.10 |
| 291 | /// regression). The retry is now modeled as this value, carried out of the |
| 292 | /// stream decoder in [`StreamOutcome`] and consumed exactly once per drop: |
| 293 | /// |
| 294 | /// * it is never persisted to the user transcript, and |
| 295 | /// * nothing it triggers is serialized into the provider request history as |
| 296 | /// a user role — the retried request is simply the persisted conversation |
| 297 | /// re-issued, ending (when a visible fragment was preserved) with that |
| 298 | /// assistant fragment so the provider continues from it. |
| 299 | #[derive(Clone, Copy, Debug, PartialEq, Eq)] |
| 300 | pub(super) enum StreamResume { |
| 301 | /// The stream died before anything actionable was streamed (#103 |
| 302 | /// Phase 3): discard the fragment and re-issue the identical request. |
| 303 | NoContentStreamDeath, |
| 304 | /// The host slept mid-stream (#2990): the partial output predates the |
| 305 | /// sleep and no operator watched it — discard and re-issue. |
| 306 | AfterSleep, |
| 307 | /// Mid-stream network drop on a headless host (v0.9.4 Terminal-Bench |
| 308 | /// P0): the fragment was never committed and no tool from it ran, so |
| 309 | /// discard and re-issue the identical request. |
| 310 | HeadlessNetworkDrop, |
| 311 | /// Mid-stream network drop in the interactive TUI. A fragment with |
| 312 | /// sendable content is preserved as the trailing assistant message and |
| 313 | /// the request is re-issued; a thinking-only fragment has nothing |
| 314 | /// visible to preserve, so it is discarded exactly like the headless |
| 315 | /// resume and the copy must never claim otherwise. |
| 316 | InteractiveNetworkDrop, |
| 317 | } |
| 318 | |
| 319 | /// Bounded authorization for drop-resume retries. |
| 320 | /// |
| 321 | /// Mechanism, not comment: [`StreamRetryBudget::authorize`] is the only way |
| 322 | /// to spend a resume and it returns `None` once `limit` resumes (default |
| 323 | /// [`MAX_STREAM_RETRIES`]) have been issued, so no call site can loop past |
| 324 | /// the budget even if a guard predicate is relaxed. A healthy stream round |
| 325 | /// resets it. |
| 326 | #[derive(Debug)] |
| 327 | pub(super) struct StreamRetryBudget { |
| 328 | spent: u32, |
| 329 | limit: u32, |
| 330 | } |
| 331 | |
| 332 | impl Default for StreamRetryBudget { |
| 333 | fn default() -> Self { |
| 334 | Self::with_limit(MAX_STREAM_RETRIES) |
| 335 | } |
| 336 | } |
| 337 | |
| 338 | impl StreamRetryBudget { |
| 339 | /// A fresh budget allowing at most `limit` resumes. |
| 340 | pub(super) fn with_limit(limit: u32) -> Self { |
| 341 | Self { spent: 0, limit } |
| 342 | } |
| 343 | |
| 344 | /// The configured resume ceiling. |
| 345 | pub(super) fn limit(&self) -> u32 { |
| 346 | self.limit |
| 347 | } |
| 348 | |
| 349 | /// Drop-resumes already issued without a healthy round in between. |
| 350 | pub(super) fn spent(&self) -> u32 { |
| 351 | self.spent |
| 352 | } |
| 353 | |
| 354 | /// Spend one resume and return its 1-based attempt number, or `None` |
| 355 | /// when the budget is exhausted. |
| 356 | pub(super) fn authorize(&mut self) -> Option<u32> { |
| 357 | if self.spent >= self.limit { |
| 358 | return None; |
| 359 | } |
| 360 | self.spent = self.spent.saturating_add(1); |
| 361 | Some(self.spent) |
| 362 | } |
| 363 | |
| 364 | /// A healthy round clears the chain: the next drop starts a fresh, |
| 365 | /// still-bounded budget. |
| 366 | pub(super) fn reset(&mut self) { |
| 367 | self.spent = 0; |
| 368 | } |
| 369 | } |
| 370 | |
| 371 | /// Wall-clock vs monotonic divergence above which we conclude the host slept |
| 372 | /// mid-stream (#2990). `Instant` pauses during system sleep (CLOCK_UPTIME_RAW |
| 373 | /// on macOS, CLOCK_MONOTONIC on Linux) while `SystemTime` keeps advancing, so |
| 374 | /// a large positive gap can only come from a suspend/resume cycle — ordinary |
| 375 | /// network flakes never produce one. Windows `Instant` may keep ticking |
| 376 | /// through sleep, in which case this simply never fires (no behavior change). |
| 377 | pub(super) const SLEEP_GAP_THRESHOLD: Duration = Duration::from_secs(10); |
| 378 | |
| 379 | /// True when the gap between wall-clock and monotonic elapsed time since the |
| 380 | /// last stream progress says the host was suspended. |
| 381 | pub(super) fn sleep_gap_detected(monotonic_elapsed: Duration, wallclock_elapsed: Duration) -> bool { |
| 382 | wallclock_elapsed.saturating_sub(monotonic_elapsed) > SLEEP_GAP_THRESHOLD |
| 383 | } |
| 384 | |
| 385 | /// Decide whether a failed stream should be silently re-issued because the |
| 386 | /// host slept mid-turn (#2990). |
| 387 | /// |
| 388 | /// Unlike the transparent retry (#103), this fires even after content has |
| 389 | /// streamed: the partial output predates the sleep, the user was not |
| 390 | /// watching, and re-running the identical request is the correct |
| 391 | /// user-visible behavior. The double-billing concern that blocks ordinary |
| 392 | /// post-content retries is accepted here because the alternative is a dead |
| 393 | /// turn the user must re-prompt (and pay for) anyway. |
| 394 | pub(super) fn should_resume_after_sleep( |
| 395 | sleep_detected: bool, |
| 396 | retry_attempts: u32, |
| 397 | retry_limit: u32, |
| 398 | cancelled: bool, |
| 399 | ) -> bool { |
| 400 | sleep_detected && retry_attempts < retry_limit && !cancelled |
| 401 | } |
| 402 | |
| 403 | /// Decide whether a failed stream should be re-issued after a mid-stream |
| 404 | /// network drop in a headless host (`exec` / stream-json / app-server), even |
| 405 | /// though content already streamed. |
| 406 | /// |
| 407 | /// This extends the #2990 sleep-resume contract to ordinary transport drops |
| 408 | /// for hosts with no operator watching: the partial assistant fragment has |
| 409 | /// not been committed to the conversation and no tool call from the |
| 410 | /// incomplete response has executed, so discarding the fragment and |
| 411 | /// re-issuing the identical request cannot duplicate side effects. The |
| 412 | /// double-billing risk that blocks post-content retries in the interactive |
| 413 | /// TUI (#103) is accepted here because the alternative is a dead turn that |
| 414 | /// forfeits the entire headless run — the exact tradeoff #2990 already makes |
| 415 | /// for sleep-resume. Interactive sessions keep the #103 surface-the-warning |
| 416 | /// behavior: the user saw the partial deltas, and replaying would render the |
| 417 | /// same prefix twice. |
| 418 | pub(super) fn should_resume_after_network_drop( |
| 419 | headless_host: bool, |
| 420 | network_class_error: bool, |
| 421 | retry_attempts: u32, |
| 422 | retry_limit: u32, |
| 423 | cancelled: bool, |
| 424 | ) -> bool { |
| 425 | headless_host && network_class_error && retry_attempts < retry_limit && !cancelled |
| 426 | } |
| 427 | |
| 428 | /// Decide whether an interactive TUI stream should be re-issued after a |
| 429 | /// mid-stream network drop, preserving a visible partial reply. |
| 430 | /// |
| 431 | /// Unlike the headless resume, this keeps a sendable fragment: the user has |
| 432 | /// already seen the deltas, so the assistant message is committed and the |
| 433 | /// re-issued request ends with that fragment, which is the provider-neutral |
| 434 | /// continuation contract. No synthetic user turn is appended — the retry is |
| 435 | /// typed state ([`StreamResume::InteractiveNetworkDrop`]), invisible to the |
| 436 | /// transcript and to the provider request history as a user role. A |
| 437 | /// thinking-only fragment preserves nothing and must not be described as a |
| 438 | /// preserved reply. Tool calls are never resumed because an incomplete tool |
| 439 | /// call could be re-issued and duplicate side effects. Bounded by |
| 440 | /// `MAX_STREAM_RETRIES` and gated on a network/timeout-class error so |
| 441 | /// model/parse/auth failures still surface normally. |
| 442 | pub(super) fn should_resume_interactive_after_network_drop( |
| 443 | terminal_chrome_enabled: bool, |
| 444 | network_class_error: bool, |
| 445 | any_content_received: bool, |
| 446 | tool_uses_empty: bool, |
| 447 | retry_attempts: u32, |
| 448 | retry_limit: u32, |
| 449 | cancelled: bool, |
| 450 | ) -> bool { |
| 451 | terminal_chrome_enabled |
| 452 | && network_class_error |
| 453 | && any_content_received |
| 454 | && tool_uses_empty |
| 455 | && retry_attempts < retry_limit |
| 456 | && !cancelled |
| 457 | } |
| 458 | |
| 459 | /// Convert low-level reqwest/hyper stream read errors into an operator-facing |
| 460 | /// message. The raw provider error remains attached, but the lead sentence |
| 461 | /// explains why Codewhale may retry before any output and why it must surface |
| 462 | /// the warning once partial output has already streamed. |
| 463 | pub(super) fn stream_read_error_user_message(message: &str, any_content_received: bool) -> String { |
| 464 | let lower = message.to_ascii_lowercase(); |
| 465 | let is_stream_read = lower.contains("stream read error") |
| 466 | || lower.contains("error decoding response body") |
| 467 | || lower.contains("chunk decode error") |
| 468 | || lower.contains("body decode"); |
| 469 | if !is_stream_read { |
| 470 | return message.to_string(); |
| 471 | } |
| 472 | |
| 473 | let retry_note = if any_content_received { |
| 474 | "Some output had already streamed, so Codewhale is surfacing the warning instead of replaying the request and risking duplicated output." |
| 475 | } else { |
| 476 | "No output had streamed yet, so Codewhale will retry automatically while retry budget remains." |
| 477 | }; |
| 478 | format!( |
| 479 | "Provider stream connection dropped while reading the response body. {retry_note} Details: {message}" |
| 480 | ) |
| 481 | } |
| 482 | |
| 483 | /// Wrapper shapes a model may emit as plain text instead of using the API tool |
| 484 | /// channel. Each pair is `(start, end)`; the tables below are projections of |
| 485 | /// this one and must stay in sync with it. |
| 486 | /// |
| 487 | /// Three families are covered: |
| 488 | /// |
| 489 | /// 1. Generic/Anthropic-style (`[TOOL_CALL]`, `<invoke …>`, `<function_calls>`). |
| 490 | /// 2. DSML wrappers, in fullwidth `|` (U+FF5C) and ASCII `|` delimiters, upper |
| 491 | /// and lower case. DeepSeek also emits a doubled-delimiter form |
| 492 | /// (`<||DSML|| calls>`) when a request offers no tools; one-shot |
| 493 | /// `codewhale exec` printed it verbatim as the answer. |
| 494 | /// 3. **DeepSeek's native tool-call tokens** (#3880). DeepSeek's chat template |
| 495 | /// separates words with `▁` (U+2581 LOWER ONE EIGHTH BLOCK), not a space or |
| 496 | /// underscore, so `<|tool▁calls▁begin|>` does not match any DSML entry and |
| 497 | /// leaked into visible output. Both the `▁` and `_` separators are listed |
| 498 | /// because a partially-normalizing tokenizer can emit either, and both |
| 499 | /// delimiter forms because the ASCII fallback shows up in some renderings. |
| 500 | /// |
| 501 | /// When adding a shape, add it here and to the two marker tables below. |
| 502 | /// `marker_tables_are_consistent` enforces that they agree. |
| 503 | pub(crate) const TOOL_CALL_MARKER_PAIRS: [(&str, &str); 30] = [ |
| 504 | ("[TOOL_CALL]", "[/TOOL_CALL]"), |
| 505 | ("<codewhale:tool_call", "</codewhale:tool_call>"), |
| 506 | ("<tool_call", "</tool_call>"), |
| 507 | ("<invoke ", "</invoke>"), |
| 508 | ("<function_calls>", "</function_calls>"), |
| 509 | ("<|DSML|tool_calls>", "</|DSML|tool_calls>"), |
| 510 | ("<|DSML|invoke ", "</|DSML|invoke>"), |
| 511 | ("<|DSML|tool_calls>", "</|DSML|tool_calls>"), |
| 512 | ("<|DSML|invoke ", "</|DSML|invoke>"), |
| 513 | ("<|dsml|tool_calls>", "</|dsml|tool_calls>"), |
| 514 | ("<|dsml|invoke ", "</|dsml|invoke>"), |
| 515 | ("<||DSML|| calls>", "</||DSML|| calls>"), |
| 516 | ("<||DSML|| invoke ", "</||DSML|| invoke>"), |
| 517 | ("<|tool_calls>", "</|tool_calls>"), |
| 518 | // DeepSeek native, fullwidth delimiters, U+2581 separator. |
| 519 | ("<|tool▁calls▁begin|>", "<|tool▁calls▁end|>"), |
| 520 | ("<|tool▁call▁begin|>", "<|tool▁call▁end|>"), |
| 521 | ("<|tool▁outputs▁begin|>", "<|tool▁outputs▁end|>"), |
| 522 | ("<|tool▁output▁begin|>", "<|tool▁output▁end|>"), |
| 523 | // DeepSeek native, ASCII delimiters, U+2581 separator. |
| 524 | ("<|tool▁calls▁begin|>", "<|tool▁calls▁end|>"), |
| 525 | ("<|tool▁call▁begin|>", "<|tool▁call▁end|>"), |
| 526 | ("<|tool▁outputs▁begin|>", "<|tool▁outputs▁end|>"), |
| 527 | ("<|tool▁output▁begin|>", "<|tool▁output▁end|>"), |
| 528 | // DeepSeek native, underscore separator. |
| 529 | ("<|tool_calls_begin|>", "<|tool_calls_end|>"), |
| 530 | ("<|tool_call_begin|>", "<|tool_call_end|>"), |
| 531 | ("<|tool_outputs_begin|>", "<|tool_outputs_end|>"), |
| 532 | ("<|tool_output_begin|>", "<|tool_output_end|>"), |
| 533 | ("<|tool_calls_begin|>", "<|tool_calls_end|>"), |
| 534 | ("<|tool_call_begin|>", "<|tool_call_end|>"), |
| 535 | ("<|tool_outputs_begin|>", "<|tool_outputs_end|>"), |
| 536 | ("<|tool_output_begin|>", "<|tool_output_end|>"), |
| 537 | ]; |
| 538 | |
| 539 | pub(crate) const TOOL_CALL_START_MARKERS: [&str; 30] = [ |
| 540 | "[TOOL_CALL]", |
| 541 | "<codewhale:tool_call", |
| 542 | "<tool_call", |
| 543 | "<invoke ", |
| 544 | "<function_calls>", |
| 545 | "<|DSML|tool_calls>", |
| 546 | "<|DSML|invoke ", |
| 547 | "<|DSML|tool_calls>", |
| 548 | "<|DSML|invoke ", |
| 549 | "<|dsml|tool_calls>", |
| 550 | "<|dsml|invoke ", |
| 551 | "<||DSML|| calls>", |
| 552 | "<||DSML|| invoke ", |
| 553 | "<|tool_calls>", |
| 554 | "<|tool▁calls▁begin|>", |
| 555 | "<|tool▁call▁begin|>", |
| 556 | "<|tool▁outputs▁begin|>", |
| 557 | "<|tool▁output▁begin|>", |
| 558 | "<|tool▁calls▁begin|>", |
| 559 | "<|tool▁call▁begin|>", |
| 560 | "<|tool▁outputs▁begin|>", |
| 561 | "<|tool▁output▁begin|>", |
| 562 | "<|tool_calls_begin|>", |
| 563 | "<|tool_call_begin|>", |
| 564 | "<|tool_outputs_begin|>", |
| 565 | "<|tool_output_begin|>", |
| 566 | "<|tool_calls_begin|>", |
| 567 | "<|tool_call_begin|>", |
| 568 | "<|tool_outputs_begin|>", |
| 569 | "<|tool_output_begin|>", |
| 570 | ]; |
| 571 | |
| 572 | pub(crate) const TOOL_CALL_END_MARKERS: [&str; 30] = [ |
| 573 | "[/TOOL_CALL]", |
| 574 | "</codewhale:tool_call>", |
| 575 | "</tool_call>", |
| 576 | "</invoke>", |
| 577 | "</function_calls>", |
| 578 | "</|DSML|tool_calls>", |
| 579 | "</|DSML|invoke>", |
| 580 | "</|DSML|tool_calls>", |
| 581 | "</|DSML|invoke>", |
| 582 | "</|dsml|tool_calls>", |
| 583 | "</|dsml|invoke>", |
| 584 | "</||DSML|| calls>", |
| 585 | "</||DSML|| invoke>", |
| 586 | "</|tool_calls>", |
| 587 | "<|tool▁calls▁end|>", |
| 588 | "<|tool▁call▁end|>", |
| 589 | "<|tool▁outputs▁end|>", |
| 590 | "<|tool▁output▁end|>", |
| 591 | "<|tool▁calls▁end|>", |
| 592 | "<|tool▁call▁end|>", |
| 593 | "<|tool▁outputs▁end|>", |
| 594 | "<|tool▁output▁end|>", |
| 595 | "<|tool_calls_end|>", |
| 596 | "<|tool_call_end|>", |
| 597 | "<|tool_outputs_end|>", |
| 598 | "<|tool_output_end|>", |
| 599 | "<|tool_calls_end|>", |
| 600 | "<|tool_call_end|>", |
| 601 | "<|tool_outputs_end|>", |
| 602 | "<|tool_output_end|>", |
| 603 | ]; |
| 604 | |
| 605 | #[derive(Debug, Default)] |
| 606 | pub(crate) struct ToolCallDeltaFilterState { |
| 607 | in_tool_call: bool, |
| 608 | marker_carry: String, |
| 609 | active_end_marker: Option<&'static str>, |
| 610 | } |
| 611 | |
| 612 | /// Compact one-shot notice emitted when a model attempts to forge a tool-call |
| 613 | /// wrapper in plain text instead of using the API tool channel. The visible |
| 614 | /// content is still scrubbed; this exists so the user can see why their text |
| 615 | /// shrank. |
| 616 | pub(crate) const FAKE_WRAPPER_NOTICE: &str = |
| 617 | "Stripped non-API tool-call wrapper from model output (use the API tool channel)"; |
| 618 | |
| 619 | /// True if `text` contains any of the known fake-wrapper start markers. Used by |
| 620 | /// the streaming loop to decide whether to emit `FAKE_WRAPPER_NOTICE`. |
| 621 | pub(crate) fn contains_fake_tool_wrapper(text: &str) -> bool { |
| 622 | TOOL_CALL_START_MARKERS.iter().any(|m| text.contains(m)) |
| 623 | } |
| 624 | |
| 625 | fn find_first_marker(text: &str, markers: &[&str]) -> Option<(usize, usize)> { |
| 626 | markers |
| 627 | .iter() |
| 628 | .filter_map(|marker| text.find(marker).map(|idx| (idx, marker.len()))) |
| 629 | .min_by_key(|(idx, _)| *idx) |
| 630 | } |
| 631 | |
| 632 | fn find_first_start_marker(text: &str) -> Option<(usize, usize, &'static str)> { |
| 633 | TOOL_CALL_MARKER_PAIRS |
| 634 | .iter() |
| 635 | .filter_map(|(start, end)| text.find(start).map(|idx| (idx, start.len(), *end))) |
| 636 | .min_by_key(|(idx, _, _)| *idx) |
| 637 | } |
| 638 | |
| 639 | /// Cheap rejection: every marker prefix ends with the marker's own first |
| 640 | /// byte, so a text without any marker's first byte cannot end with one. |
| 641 | /// Every tool-call marker starts with `<` or `[`, so this is one scan of a |
| 642 | /// usually short delta and keeps per-token `ends_with` probing off the hot |
| 643 | /// path for plain prose deltas. |
| 644 | fn fast_reject_marker_text(text: &str) -> bool { |
| 645 | !text.bytes().any(|b| b == b'<' || b == b'[') |
| 646 | } |
| 647 | |
| 648 | fn trailing_marker_prefix_len(text: &str, markers: &[&str]) -> usize { |
| 649 | if fast_reject_marker_text(text) { |
| 650 | return 0; |
| 651 | } |
| 652 | markers |
| 653 | .iter() |
| 654 | .flat_map(|marker| { |
| 655 | marker |
| 656 | .char_indices() |
| 657 | .map(|(idx, _)| idx) |
| 658 | .filter(|idx| *idx > 0) |
| 659 | .chain(std::iter::once(marker.len())) |
| 660 | .filter(|idx| *idx < marker.len()) |
| 661 | .filter(|idx| { |
| 662 | let prefix = &marker[..*idx]; |
| 663 | text.ends_with(prefix) |
| 664 | }) |
| 665 | }) |
| 666 | .max() |
| 667 | .unwrap_or(0) |
| 668 | } |
| 669 | |
| 670 | fn trailing_start_marker_prefix_len(text: &str) -> usize { |
| 671 | if fast_reject_marker_text(text) { |
| 672 | return 0; |
| 673 | } |
| 674 | TOOL_CALL_MARKER_PAIRS |
| 675 | .iter() |
| 676 | .flat_map(|(marker, _)| { |
| 677 | marker |
| 678 | .char_indices() |
| 679 | .map(|(idx, _)| idx) |
| 680 | .filter(|idx| *idx > 0) |
| 681 | .chain(std::iter::once(marker.len())) |
| 682 | .filter(|idx| *idx < marker.len()) |
| 683 | .filter(|idx| { |
| 684 | let prefix = &marker[..*idx]; |
| 685 | text.ends_with(prefix) |
| 686 | }) |
| 687 | }) |
| 688 | .max() |
| 689 | .unwrap_or(0) |
| 690 | } |
| 691 | |
| 692 | #[cfg(test)] |
| 693 | pub(crate) fn filter_tool_call_delta(delta: &str, in_tool_call: &mut bool) -> String { |
| 694 | let mut state = ToolCallDeltaFilterState { |
| 695 | in_tool_call: *in_tool_call, |
| 696 | ..ToolCallDeltaFilterState::default() |
| 697 | }; |
| 698 | let output = filter_tool_call_delta_with_state(delta, &mut state); |
| 699 | *in_tool_call = state.in_tool_call; |
| 700 | output |
| 701 | } |
| 702 | |
| 703 | pub(crate) fn filter_tool_call_delta_with_state( |
| 704 | delta: &str, |
| 705 | state: &mut ToolCallDeltaFilterState, |
| 706 | ) -> String { |
| 707 | if delta.is_empty() { |
| 708 | return String::new(); |
| 709 | } |
| 710 | |
| 711 | let chunk; |
| 712 | let mut rest = if state.marker_carry.is_empty() { |
| 713 | delta |
| 714 | } else { |
| 715 | chunk = format!("{}{delta}", state.marker_carry); |
| 716 | state.marker_carry.clear(); |
| 717 | &chunk |
| 718 | }; |
| 719 | let mut output = String::new(); |
| 720 | |
| 721 | loop { |
| 722 | if state.in_tool_call { |
| 723 | // Close only on the opener's own end marker when it is known. |
| 724 | // Falling back to any end marker let a nested closer inside a |
| 725 | // DSML block (`</|DSML|invoke>`) end the block whenever a delta |
| 726 | // lacked the outer closer, leaking `</|DSML|tool_calls>`. |
| 727 | let active_end_marker = state.active_end_marker; |
| 728 | let found = match active_end_marker { |
| 729 | Some(marker) => rest.find(marker).map(|idx| (idx, marker.len())), |
| 730 | None => find_first_marker(rest, &TOOL_CALL_END_MARKERS), |
| 731 | }; |
| 732 | let Some((idx, len)) = found else { |
| 733 | let keep = active_end_marker.map_or_else( |
| 734 | || trailing_marker_prefix_len(rest, &TOOL_CALL_END_MARKERS), |
| 735 | |marker| trailing_marker_prefix_len(rest, &[marker]), |
| 736 | ); |
| 737 | if keep > 0 { |
| 738 | state.marker_carry.push_str(&rest[rest.len() - keep..]); |
| 739 | } |
| 740 | break; |
| 741 | }; |
| 742 | rest = &rest[idx + len..]; |
| 743 | state.in_tool_call = false; |
| 744 | state.active_end_marker = None; |
| 745 | } else { |
| 746 | let Some((idx, len, end_marker)) = find_first_start_marker(rest) else { |
| 747 | let keep = trailing_start_marker_prefix_len(rest); |
| 748 | if keep > 0 { |
| 749 | let split = rest.len() - keep; |
| 750 | output.push_str(&rest[..split]); |
| 751 | state.marker_carry.push_str(&rest[split..]); |
| 752 | } else { |
| 753 | output.push_str(rest); |
| 754 | } |
| 755 | break; |
| 756 | }; |
| 757 | output.push_str(&rest[..idx]); |
| 758 | rest = &rest[idx + len..]; |
| 759 | state.in_tool_call = true; |
| 760 | state.active_end_marker = Some(end_marker); |
| 761 | } |
| 762 | } |
| 763 | |
| 764 | output |
| 765 | } |
| 766 | |
| 767 | pub(crate) fn flush_tool_call_delta_state(state: &mut ToolCallDeltaFilterState) -> String { |
| 768 | if state.in_tool_call { |
| 769 | state.marker_carry.clear(); |
| 770 | return String::new(); |
| 771 | } |
| 772 | std::mem::take(&mut state.marker_carry) |
| 773 | } |
| 774 |