返回 CodeWhale
streaming.rs
根目录 / crates / tui / src / core / engine / streaming.rs
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
774 lines RUST