返回 CodeWhale
stream_decoder.rs
根目录 / crates / tui / src / client / chat / tests / stream_decoder.rs
1 //! Drive `parse_sse_chunk` (the in-place SSE event extractor) over canned
2 //! chunk sequences. The full `handle_chat_completion_stream` path needs a
3 //! live `reqwest::Response` so it isn't unit-testable without a mock HTTP
4 //! harness (issue #69 tracks that). For #103 we exercise the chunk decoder
5 //! directly to verify each "class of stream failure" the engine relies on.
6 use super::*;
7 use crate::client::wire::{SseLineDecoder, SseLineError};
8 use codewhale_models::{ContentBlockStart, Delta, StreamEvent};
9
10 /// Decode a raw SSE-data JSON chunk into our internal events, mirroring
11 /// the per-event call shape used by `handle_chat_completion_stream`.
12 fn decode_chunk(json_text: &str) -> Vec<StreamEvent> {
13 decode_chunk_with_reasoning(json_text, true)
14 }
15
16 fn decode_chunk_with_reasoning(json_text: &str, is_reasoning_model: bool) -> Vec<StreamEvent> {
17 let chunk: Value = serde_json::from_str(json_text).expect("valid SSE JSON");
18 let mut content_index = 0u32;
19 let mut text_started = false;
20 let mut thinking_started = false;
21 let mut tool_indices = std::collections::HashMap::new();
22 let mut reasoning_detail_buffers = std::collections::HashMap::new();
23 parse_sse_chunk(
24 &chunk,
25 &mut content_index,
26 &mut text_started,
27 &mut thinking_started,
28 &mut tool_indices,
29 &mut reasoning_detail_buffers,
30 is_reasoning_model,
31 )
32 }
33
34 fn decode_chunks_with_style(
35 chunks: &[&str],
36 reasoning_stream_style: ReasoningStreamStyle,
37 ) -> Vec<StreamEvent> {
38 let mut content_index = 0u32;
39 let mut text_started = false;
40 let mut thinking_started = false;
41 let mut tool_indices = std::collections::HashMap::new();
42 let mut reasoning_detail_buffers = std::collections::HashMap::new();
43 let mut inline_reasoning_tags = InlineReasoningTagState::default();
44 let mut events = Vec::new();
45
46 for chunk in chunks {
47 let value: Value = serde_json::from_str(chunk).expect("valid SSE JSON");
48 events.extend(parse_sse_chunk_with_reasoning_style(
49 &value,
50 &mut content_index,
51 &mut text_started,
52 &mut thinking_started,
53 &mut tool_indices,
54 &mut reasoning_detail_buffers,
55 &mut inline_reasoning_tags,
56 reasoning_stream_style,
57 ));
58 }
59 events
60 }
61
62 /// Drive the Chat Completions SSE path with raw byte chunks so tests can
63 /// split a multi-byte UTF-8 character across HTTP/2-style DATA boundaries.
64 fn decode_sse_byte_chunks(chunks: &[&[u8]]) -> Result<Vec<StreamEvent>, SseLineError> {
65 struct FrameState {
66 line_buf: String,
67 content_index: u32,
68 text_started: bool,
69 thinking_started: bool,
70 tool_indices: std::collections::HashMap<u32, u32>,
71 reasoning_detail_buffers: std::collections::HashMap<u32, String>,
72 inline_reasoning_tags: InlineReasoningTagState,
73 events: Vec<StreamEvent>,
74 }
75
76 impl FrameState {
77 fn new() -> Self {
78 Self {
79 line_buf: String::new(),
80 content_index: 0,
81 text_started: false,
82 thinking_started: false,
83 tool_indices: std::collections::HashMap::new(),
84 reasoning_detail_buffers: std::collections::HashMap::new(),
85 inline_reasoning_tags: InlineReasoningTagState::default(),
86 events: Vec::new(),
87 }
88 }
89
90 fn handle_line(&mut self, line: &str) -> bool {
91 if line.is_empty() {
92 return matches!(self.flush_frame(), SseDataFrame::Done);
93 }
94 if let Some(data) = extract_sse_data_value(line) {
95 if !self.line_buf.is_empty() {
96 self.line_buf.push('\n');
97 }
98 self.line_buf.push_str(data);
99 }
100 false
101 }
102
103 fn flush_frame(&mut self) -> SseDataFrame {
104 if self.line_buf.is_empty() {
105 return SseDataFrame::Events(Vec::new());
106 }
107 let data = std::mem::take(&mut self.line_buf);
108 match parse_sse_data_frame(
109 &data,
110 &mut self.content_index,
111 &mut self.text_started,
112 &mut self.thinking_started,
113 &mut self.tool_indices,
114 &mut self.reasoning_detail_buffers,
115 &mut self.inline_reasoning_tags,
116 ReasoningStreamStyle::SeparateField,
117 ) {
118 SseDataFrame::Done => SseDataFrame::Done,
119 SseDataFrame::Events(frame_events) => {
120 self.events.extend(frame_events);
121 SseDataFrame::Events(Vec::new())
122 }
123 }
124 }
125 }
126
127 let mut decoder = SseLineDecoder::new();
128 let mut state = FrameState::new();
129 for chunk in chunks {
130 for line in decoder.push(chunk)? {
131 if state.handle_line(&line) {
132 return Ok(state.events);
133 }
134 }
135 }
136 if let Some(line) = decoder.finish()?
137 && state.handle_line(&line)
138 {
139 return Ok(state.events);
140 }
141 state.flush_frame();
142 Ok(state.events)
143 }
144
145 fn cjk_content_sse(text: &str) -> Vec<u8> {
146 let payload = serde_json::json!({
147 "choices": [{ "delta": { "content": text } }]
148 });
149 format!("data: {payload}\n\n").into_bytes()
150 }
151
152 fn mid_char_split(bytes: &[u8], ch: char) -> usize {
153 let needle = ch.to_string();
154 let start = bytes
155 .windows(needle.len())
156 .position(|window| window == needle.as_bytes())
157 .unwrap_or_else(|| panic!("{ch:?} present in SSE frame"));
158 start + 1
159 }
160
161 fn text_delta_text(events: &[StreamEvent]) -> String {
162 events
163 .iter()
164 .filter_map(|event| match event {
165 StreamEvent::ContentBlockDelta {
166 delta: Delta::TextDelta { text },
167 ..
168 } => Some(text.as_str()),
169 _ => None,
170 })
171 .collect()
172 }
173
174 fn thinking_delta_text(events: &[StreamEvent]) -> String {
175 events
176 .iter()
177 .filter_map(|event| match event {
178 StreamEvent::ContentBlockDelta {
179 delta: Delta::ThinkingDelta { thinking },
180 ..
181 } => Some(thinking.as_str()),
182 _ => None,
183 })
184 .collect()
185 }
186
187 #[test]
188 fn decoder_reassembles_cjk_split_across_byte_chunks() {
189 let frame = cjk_content_sse("你好世界");
190 let split = mid_char_split(&frame, '好');
191 let events = decode_sse_byte_chunks(&[&frame[..split], &frame[split..]]).expect("valid utf-8");
192 let text = text_delta_text(&events);
193 assert_eq!(text, "你好世界");
194 assert!(
195 !text.contains('\u{FFFD}'),
196 "HTTP/2 mid-character split must not substitute U+FFFD; got {text:?}"
197 );
198 }
199
200 #[test]
201 fn decoder_reassembles_emoji_and_cjk_fed_one_byte_at_a_time() {
202 let frame = cjk_content_sse("你好🌊世界");
203 let chunks: Vec<&[u8]> = frame.chunks(1).collect();
204 let events = decode_sse_byte_chunks(&chunks).expect("valid utf-8");
205 let text = text_delta_text(&events);
206 assert_eq!(text, "你好🌊世界");
207 assert!(
208 !text.contains('\u{FFFD}'),
209 "byte-at-a-time feed garbled: {text:?}"
210 );
211 }
212
213 #[test]
214 fn decoder_rejects_invalid_sse_bytes_without_replacement() {
215 let mut frame = cjk_content_sse("ok");
216 // Bare 0xFF is never valid UTF-8. Insert it inside the first SSE line.
217 let newline = frame.iter().position(|&b| b == b'\n').expect("SSE line");
218 frame.insert(newline, 0xFF);
219 let result = decode_sse_byte_chunks(&[&frame]);
220 let err = result.expect_err("invalid SSE bytes must fail closed");
221 assert!(
222 !err.to_string().contains('\u{FFFD}'),
223 "error path must not substitute U+FFFD: {err}"
224 );
225
226 // Unterminated tail of continuation bytes: fail on flush, no replacement.
227 let result = decode_sse_byte_chunks(&[&[0x80, 0xBF]]);
228 let err = result.expect_err("invalid unterminated flush must fail closed");
229 assert!(!err.to_string().contains('\u{FFFD}'));
230 }
231
232 #[test]
233 fn decoder_emits_text_delta_for_content_chunk() {
234 // The "happy" first chunk: a normal content delta. The engine treats
235 // this as `any_content_received = true` and would NOT transparently
236 // retry on a subsequent error.
237 let events = decode_chunk(r#"{"choices":[{"delta":{"content":"hello"}}]}"#);
238 assert!(
239 matches!(
240 events.first(),
241 Some(StreamEvent::ContentBlockStart {
242 content_block: ContentBlockStart::Text { .. },
243 ..
244 })
245 ),
246 "first event should open a text block; got {events:?}"
247 );
248 assert!(
249 events
250 .iter()
251 .any(|e| matches!(e, StreamEvent::ContentBlockDelta {
252 delta: Delta::TextDelta { text },
253 ..
254 } if text == "hello")),
255 "should yield a TextDelta carrying 'hello'; got {events:?}"
256 );
257 }
258
259 #[test]
260 fn decoder_emits_thinking_delta_for_reasoning_chunk() {
261 // V4 thinking models surface reasoning_content first — the engine
262 // also counts these as content received (so a subsequent stream error
263 // surfaces rather than retrying transparently).
264 let events = decode_chunk(r#"{"choices":[{"delta":{"reasoning_content":"plan..."}}]}"#);
265 assert!(
266 matches!(
267 events.first(),
268 Some(StreamEvent::ContentBlockStart {
269 content_block: ContentBlockStart::Thinking { .. },
270 ..
271 })
272 ),
273 "first event should open a thinking block; got {events:?}"
274 );
275 assert!(
276 events
277 .iter()
278 .any(|e| matches!(e, StreamEvent::ContentBlockDelta {
279 delta: Delta::ThinkingDelta { thinking },
280 ..
281 } if thinking == "plan...")),
282 "should yield a ThinkingDelta carrying 'plan...'; got {events:?}"
283 );
284 }
285
286 #[test]
287 fn decoder_streams_moonshot_multi_chunk_reasoning_as_thinking() {
288 // #3016: recorded shape from Moonshot's native endpoint — kimi-k2.6
289 // streams `reasoning_content` deltas before the answer text. The
290 // thinking deltas must accumulate into ONE thinking block and the
291 // answer must arrive as text, not be glued into the trace.
292 let chunks = [
293 r#"{"id":"cmpl-kimi","model":"kimi-k2.6","choices":[{"index":0,"delta":{"role":"assistant","reasoning_content":"Let me check"}}]}"#,
294 r#"{"id":"cmpl-kimi","model":"kimi-k2.6","choices":[{"index":0,"delta":{"reasoning_content":" the config."}}]}"#,
295 r#"{"id":"cmpl-kimi","model":"kimi-k2.6","choices":[{"index":0,"delta":{"content":"The answer is 42."}}]}"#,
296 ];
297
298 let is_reasoning =
299 is_reasoning_model_for_stream(crate::config::ProviderKind::Moonshot, "kimi-k2.6");
300 let mut content_index = 0u32;
301 let mut text_started = false;
302 let mut thinking_started = false;
303 let mut tool_indices = std::collections::HashMap::new();
304 let mut reasoning_detail_buffers = std::collections::HashMap::new();
305 let mut events = Vec::new();
306 for chunk in chunks {
307 let value: Value = serde_json::from_str(chunk).expect("valid SSE JSON");
308 events.extend(parse_sse_chunk(
309 &value,
310 &mut content_index,
311 &mut text_started,
312 &mut thinking_started,
313 &mut tool_indices,
314 &mut reasoning_detail_buffers,
315 is_reasoning,
316 ));
317 }
318
319 let thinking: String = events
320 .iter()
321 .filter_map(|event| match event {
322 StreamEvent::ContentBlockDelta {
323 delta: Delta::ThinkingDelta { thinking },
324 ..
325 } => Some(thinking.as_str()),
326 _ => None,
327 })
328 .collect();
329 assert_eq!(thinking, "Let me check the config.");
330
331 let thinking_starts = events
332 .iter()
333 .filter(|event| {
334 matches!(
335 event,
336 StreamEvent::ContentBlockStart {
337 content_block: ContentBlockStart::Thinking { .. },
338 ..
339 }
340 )
341 })
342 .count();
343 assert_eq!(thinking_starts, 1, "one thinking block: {events:?}");
344
345 let text: String = events
346 .iter()
347 .filter_map(|event| match event {
348 StreamEvent::ContentBlockDelta {
349 delta: Delta::TextDelta { text },
350 ..
351 } => Some(text.as_str()),
352 _ => None,
353 })
354 .collect();
355 assert_eq!(text, "The answer is 42.");
356 }
357
358 #[test]
359 fn decoder_accepts_openrouter_reasoning_delta_with_extra_fields() {
360 let events = decode_chunk(
361 r#"{"id":"or-1","choices":[{"delta":{"reasoning":"openrouter thought","reasoning_details":[{"type":"summary","text":"extra"}],"native_finish_reason":null}}],"usage":{"completion_tokens_details":{"reasoning_tokens":3}}}"#,
362 );
363
364 assert!(
365 events.iter().any(|e| matches!(
366 e,
367 StreamEvent::ContentBlockDelta {
368 delta: Delta::ThinkingDelta { thinking },
369 ..
370 } if thinking == "openrouter thought"
371 )),
372 "OpenRouter-style reasoning deltas with extra fields should not crash decoding; got {events:?}"
373 );
374 }
375
376 #[test]
377 fn decoder_streams_minimax_reasoning_details_as_incremental_thinking() {
378 // MiniMax's reasoning_split stream reports reasoning_details text as
379 // a cumulative buffer. Emit only the suffix so the Thinking cell does
380 // not duplicate earlier reasoning chunks.
381 let chunks = [
382 r#"{"id":"minimax-1","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"text","text":"Inspect"}]}}]}"#,
383 r#"{"id":"minimax-1","choices":[{"index":0,"delta":{"reasoning_details":[{"type":"text","text":"Inspect config"}]}}]}"#,
384 r#"{"id":"minimax-1","choices":[{"index":0,"delta":{"content":"Done."}}]}"#,
385 ];
386
387 let is_reasoning = is_reasoning_model_for_stream(ProviderKind::Minimax, "MiniMax-M3");
388 let mut content_index = 0u32;
389 let mut text_started = false;
390 let mut thinking_started = false;
391 let mut tool_indices = std::collections::HashMap::new();
392 let mut reasoning_detail_buffers = std::collections::HashMap::new();
393 let mut events = Vec::new();
394 for chunk in chunks {
395 let value: Value = serde_json::from_str(chunk).expect("valid SSE JSON");
396 events.extend(parse_sse_chunk(
397 &value,
398 &mut content_index,
399 &mut text_started,
400 &mut thinking_started,
401 &mut tool_indices,
402 &mut reasoning_detail_buffers,
403 is_reasoning,
404 ));
405 }
406
407 let thinking: String = events
408 .iter()
409 .filter_map(|event| match event {
410 StreamEvent::ContentBlockDelta {
411 delta: Delta::ThinkingDelta { thinking },
412 ..
413 } => Some(thinking.as_str()),
414 _ => None,
415 })
416 .collect();
417 assert_eq!(thinking, "Inspect config");
418
419 assert!(!events.iter().any(|event| matches!(
420 event,
421 StreamEvent::ContentBlockDelta {
422 delta: Delta::TextDelta { text },
423 ..
424 } if text == "Inspect" || text == "Inspect config"
425 )));
426 }
427
428 #[test]
429 fn modelstudio_streams_reasoning_content_as_thinking() {
430 // Recorded-style DashScope OpenAI-compatible frames (shape lifted from
431 // Model Studio's deep-thinking docs): reasoning streams in
432 // `delta.reasoning_content`, the answer in `delta.content`, and a
433 // trailing usage-only chunk closes the stream.
434 let chunks = [
435 r#"{"choices":[{"delta":{"content":null,"role":"assistant","reasoning_content":""},"index":0,"logprobs":null,"finish_reason":null}],"object":"chat.completion.chunk","usage":null,"model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
436 r#"{"choices":[{"delta":{"reasoning_content":"Let me think"},"index":0}],"object":"chat.completion.chunk","model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
437 r#"{"choices":[{"delta":{"reasoning_content":" about this."},"index":0}],"object":"chat.completion.chunk","model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
438 r#"{"choices":[{"delta":{"content":"The answer."},"index":0}],"object":"chat.completion.chunk","model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
439 r#"{"choices":[{"finish_reason":"stop","delta":{"content":"","reasoning_content":null},"index":0}],"object":"chat.completion.chunk","model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
440 r#"{"choices":[],"object":"chat.completion.chunk","usage":{"prompt_tokens":10,"completion_tokens":30,"total_tokens":40,"completion_tokens_details":{"reasoning_tokens":20}},"model":"qwen3.8-max","id":"chatcmpl-ms-1"}"#,
441 ];
442
443 // Both OpenAI-dialect plans classify their reasoning catalog.
444 for (provider, base_url, model) in [
445 (
446 ProviderKind::ModelstudioTokenPlan,
447 crate::config::DEFAULT_MODELSTUDIO_TOKEN_PLAN_BASE_URL,
448 "qwen3.8-max",
449 ),
450 (
451 ProviderKind::ModelstudioCodingPlan,
452 crate::config::DEFAULT_MODELSTUDIO_CODING_PLAN_BASE_URL,
453 "qwen3.7-plus",
454 ),
455 ] {
456 let style = reasoning_stream_style_for_route(provider, base_url, model, None);
457 assert_eq!(style, ReasoningStreamStyle::SeparateField, "{provider:?}");
458
459 let mut content_index = 0u32;
460 let mut text_started = false;
461 let mut thinking_started = false;
462 let mut tool_indices = std::collections::HashMap::new();
463 let mut reasoning_detail_buffers = std::collections::HashMap::new();
464 let mut inline_reasoning_tags = InlineReasoningTagState::default();
465 let mut events = Vec::new();
466 for chunk in chunks {
467 let value: Value = serde_json::from_str(chunk).expect("valid SSE JSON");
468 events.extend(parse_sse_chunk_with_reasoning_style(
469 &value,
470 &mut content_index,
471 &mut text_started,
472 &mut thinking_started,
473 &mut tool_indices,
474 &mut reasoning_detail_buffers,
475 &mut inline_reasoning_tags,
476 style,
477 ));
478 }
479
480 let thinking: String = events
481 .iter()
482 .filter_map(|event| match event {
483 StreamEvent::ContentBlockDelta {
484 delta: Delta::ThinkingDelta { thinking },
485 ..
486 } => Some(thinking.as_str()),
487 _ => None,
488 })
489 .collect();
490 assert_eq!(thinking, "Let me think about this.", "{provider:?}");
491
492 let text: String = events
493 .iter()
494 .filter_map(|event| match event {
495 StreamEvent::ContentBlockDelta {
496 delta: Delta::TextDelta { text },
497 ..
498 } => Some(text.as_str()),
499 _ => None,
500 })
501 .collect();
502 assert_eq!(text, "The answer.", "{provider:?}");
503
504 // The trailing usage chunk still surfaces token accounting.
505 assert!(
506 events.iter().any(|event| matches!(
507 event,
508 StreamEvent::MessageDelta { usage: Some(usage), .. }
509 if usage.output_tokens == 30
510 )),
511 "{provider:?}: {events:?}"
512 );
513 }
514
515 // A model id the catalog does not know still routes the reasoning field
516 // to Thinking (#6501); a model that sends no reasoning field gets none.
517 let style = reasoning_stream_style_for_route(
518 ProviderKind::ModelstudioTokenPlan,
519 crate::config::DEFAULT_MODELSTUDIO_TOKEN_PLAN_BASE_URL,
520 "qwen3.8-max-lite-unknown",
521 None,
522 );
523 assert_eq!(style, ReasoningStreamStyle::SeparateField);
524 }
525
526 #[test]
527 fn decoder_does_not_render_reasoning_as_text_for_known_provider_models() {
528 let mut content_index = 0u32;
529 let mut text_started = false;
530 let mut thinking_started = false;
531 let mut tool_indices = std::collections::HashMap::new();
532 let mut reasoning_detail_buffers = std::collections::HashMap::new();
533 let is_reasoning_model =
534 is_reasoning_model_for_stream(ProviderKind::XiaomiMimo, "mimo-v2.5-pro");
535 let events = parse_sse_chunk(
536 &serde_json::json!({
537 "choices": [{
538 "delta": {
539 "reasoning_content": "private plan"
540 }
541 }]
542 }),
543 &mut content_index,
544 &mut text_started,
545 &mut thinking_started,
546 &mut tool_indices,
547 &mut reasoning_detail_buffers,
548 is_reasoning_model,
549 );
550
551 assert!(events.iter().any(|event| matches!(
552 event,
553 StreamEvent::ContentBlockDelta {
554 delta: Delta::ThinkingDelta { thinking },
555 ..
556 } if thinking == "private plan"
557 )));
558 assert!(!events.iter().any(|event| matches!(
559 event,
560 StreamEvent::ContentBlockDelta {
561 delta: Delta::TextDelta { text },
562 ..
563 } if text == "private plan"
564 )));
565 }
566
567 #[test]
568 fn decoder_treats_reasoning_content_as_text_when_provider_does_not_support_reasoning() {
569 let events = decode_chunk_with_reasoning(
570 r#"{"choices":[{"delta":{"reasoning_content":"hello"}}]}"#,
571 false,
572 );
573
574 assert!(
575 matches!(
576 events.first(),
577 Some(StreamEvent::ContentBlockStart {
578 content_block: ContentBlockStart::Text { .. },
579 ..
580 })
581 ),
582 "first event should open a text block; got {events:?}"
583 );
584 assert!(
585 events.iter().any(|e| matches!(
586 e,
587 StreamEvent::ContentBlockDelta {
588 delta: Delta::TextDelta { text },
589 ..
590 } if text == "hello"
591 )),
592 "should yield a TextDelta carrying 'hello'; got {events:?}"
593 );
594 assert!(
595 !events.iter().any(|e| matches!(
596 e,
597 StreamEvent::ContentBlockDelta {
598 delta: Delta::ThinkingDelta { .. },
599 ..
600 }
601 )),
602 "should not emit thinking deltas for generic providers; got {events:?}"
603 );
604 }
605
606 #[test]
607 fn reasoning_style_separate_field_routes_reasoning_to_thinking() {
608 let events = decode_chunks_with_style(
609 &[
610 r#"{"choices":[{"delta":{"reasoning_content":"private plan"}}]}"#,
611 r#"{"choices":[{"delta":{"content":"Public answer."}}]}"#,
612 ],
613 ReasoningStreamStyle::SeparateField,
614 );
615
616 assert_eq!(thinking_delta_text(&events), "private plan");
617 assert_eq!(text_delta_text(&events), "Public answer.");
618 }
619
620 #[test]
621 fn exact_kimi_code_k3_streams_reasoning_content_as_thinking() {
622 let style = reasoning_stream_style_for_route(
623 ProviderKind::Moonshot,
624 crate::config::DEFAULT_KIMI_CODE_BASE_URL,
625 crate::config::KIMI_CODE_K3_MODEL,
626 None,
627 );
628 assert_eq!(style, ReasoningStreamStyle::SeparateField);
629
630 let events = decode_chunks_with_style(
631 &[r#"{"choices":[{"delta":{"reasoning_content":"private K3 plan"}}]}"#],
632 style,
633 );
634 assert_eq!(thinking_delta_text(&events), "private K3 plan");
635 assert_eq!(text_delta_text(&events), "");
636
637 let generic_style = reasoning_stream_style_for_route(
638 ProviderKind::Moonshot,
639 crate::config::DEFAULT_MOONSHOT_BASE_URL,
640 crate::config::KIMI_CODE_K3_MODEL,
641 None,
642 );
643 assert_eq!(generic_style, ReasoningStreamStyle::SeparateField);
644 }
645
646 /// #6501 regression: the founder's grok-4.7 (xAI) and mimo-v2.6-pro
647 /// (Xiaomi MiMo) sessions persisted the reasoning summary glued to the answer
648 /// in one Text block ("...I should help them find large f...Sure, I'd be
649 /// happy to help"). Decode that stream shape through the real route style.
650 #[test]
651 fn issue_6501_reasoning_field_never_leaks_into_answer_text_on_unlisted_routes() {
652 for (provider, base_url, model) in [
653 (
654 ProviderKind::Xai,
655 crate::config::DEFAULT_XAI_BASE_URL,
656 "grok-4.7",
657 ),
658 (
659 ProviderKind::XiaomiMimo,
660 "https://token-plan-sgp.xiaomimimo.com/v1",
661 "mimo-v2.7-pro-unreleased",
662 ),
663 (
664 ProviderKind::Openai,
665 "https://gateway.example.test/v1",
666 "some-new-reasoner",
667 ),
668 ] {
669 let style = reasoning_stream_style_for_route(provider, base_url, model, None);
670 assert_eq!(style, ReasoningStreamStyle::SeparateField, "{provider:?}");
671 let events = decode_chunks_with_style(
672 &[
673 r#"{"choices":[{"delta":{"role":"assistant","reasoning_content":"The user wants help finding large files"}}]}"#,
674 r#"{"choices":[{"delta":{"reasoning_content":"..."}}]}"#,
675 r#"{"choices":[{"delta":{"content":"Sure, I'd be happy to help."}}]}"#,
676 r#"{"choices":[{"delta":{"reasoning":"then a tool"}}]}"#,
677 r#"{"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","type":"function","function":{"name":"exec_shell","arguments":"{}"}}]}}]}"#,
678 r#"{"choices":[{"finish_reason":"tool_calls"}]}"#,
679 ],
680 style,
681 );
682 assert_eq!(
683 thinking_delta_text(&events),
684 "The user wants help finding large files...then a tool",
685 "{provider:?}"
686 );
687 assert_eq!(
688 text_delta_text(&events),
689 "Sure, I'd be happy to help.",
690 "{provider:?}: reasoning must not reach answer text"
691 );
692 assert!(
693 events.iter().any(|event| matches!(
694 event,
695 StreamEvent::ContentBlockStart {
696 content_block: ContentBlockStart::ToolUse { .. },
697 ..
698 }
699 )),
700 "{provider:?}: reasoning -> tool call transition must keep the tool call"
701 );
702 }
703
704 // The explicit opt-out keeps the legacy pass-through for gateways that
705 // really stream their answer in `reasoning_content`.
706 let passthrough = decode_chunks_with_style(
707 &[r#"{"choices":[{"delta":{"reasoning_content":"answer via reasoning field"}}]}"#],
708 reasoning_stream_style_for_route(
709 ProviderKind::Openai,
710 "https://gateway.example.test/v1",
711 "some-new-reasoner",
712 Some("none"),
713 ),
714 );
715 assert_eq!(thinking_delta_text(&passthrough), "");
716 assert_eq!(text_delta_text(&passthrough), "answer via reasoning field");
717 }
718
719 #[test]
720 fn reasoning_style_inline_tags_routes_think_blocks_to_thinking() {
721 let events = decode_chunks_with_style(
722 &[
723 r#"{"choices":[{"delta":{"content":"Before <thi"}}]}"#,
724 r#"{"choices":[{"delta":{"content":"nk>private plan</thi"}}]}"#,
725 r#"{"choices":[{"delta":{"content":"nk> after."}}]}"#,
726 ],
727 ReasoningStreamStyle::InlineTags,
728 );
729
730 assert_eq!(thinking_delta_text(&events), "private plan");
731 assert_eq!(text_delta_text(&events), "Before after.");
732 assert!(
733 !text_delta_text(&events).contains("<think>"),
734 "inline reasoning tags must not leak into visible text: {events:?}"
735 );
736 }
737
738 #[test]
739 fn reasoning_style_inline_tags_flushes_unclosed_think_at_stream_end() {
740 let events = decode_chunks_with_style(
741 &[
742 r#"{"choices":[{"delta":{"content":"Before <think>partial reasoning"}}]}"#,
743 r#"{"choices":[{"finish_reason":"stop"}]}"#,
744 ],
745 ReasoningStreamStyle::InlineTags,
746 );
747
748 assert_eq!(thinking_delta_text(&events), "partial reasoning");
749 assert_eq!(text_delta_text(&events), "Before ");
750 }
751
752 #[test]
753 fn reasoning_style_inline_tags_ignores_separate_reasoning_field() {
754 let events = decode_chunks_with_style(
755 &[
756 r#"{"choices":[{"delta":{"reasoning_content":"metadata","content":"<think>tagged</think> answer"}}]}"#,
757 ],
758 ReasoningStreamStyle::InlineTags,
759 );
760
761 assert_eq!(thinking_delta_text(&events), "tagged");
762 assert_eq!(text_delta_text(&events), " answer");
763 }
764
765 #[test]
766 fn reasoning_style_none_keeps_inline_tags_visible_text() {
767 let events = decode_chunks_with_style(
768 &[r#"{"choices":[{"delta":{"content":"<think>visible</think> answer"}}]}"#],
769 ReasoningStreamStyle::None,
770 );
771
772 assert_eq!(thinking_delta_text(&events), "");
773 assert_eq!(text_delta_text(&events), "<think>visible</think> answer");
774 }
775
776 #[test]
777 fn configured_reasoning_style_overrides_route_default() {
778 assert_eq!(
779 reasoning_stream_style_for_stream(ProviderKind::Openai, "custom-minimax", None),
780 ReasoningStreamStyle::SeparateField
781 );
782 assert_eq!(
783 reasoning_stream_style_for_stream(
784 ProviderKind::Openai,
785 "custom-minimax",
786 Some("inline-tags")
787 ),
788 ReasoningStreamStyle::InlineTags
789 );
790 assert_eq!(
791 reasoning_stream_style_for_stream(ProviderKind::XiaomiMimo, "mimo-v2.5-pro", None),
792 ReasoningStreamStyle::SeparateField
793 );
794 assert_eq!(
795 reasoning_stream_style_for_stream(ProviderKind::XiaomiMimo, "mimo-v2.5-pro", Some("none")),
796 ReasoningStreamStyle::None
797 );
798 }
799
800 #[test]
801 fn decoder_yields_no_events_for_keepalive_chunk() {
802 // DeepSeek often sends `{"choices":[]}` keepalive chunks before
803 // emitting real content. The engine MUST treat a stream error after
804 // these as "no content received" and be eligible for transparent
805 // retry — assert here that the decoder yields no payload events.
806 let events = decode_chunk(r#"{"choices":[]}"#);
807 assert!(
808 events.is_empty(),
809 "empty-choices chunk must produce no events; got {events:?}"
810 );
811 }
812
813 #[test]
814 fn decoder_treats_done_frame_as_terminal() {
815 let mut content_index = 0u32;
816 let mut text_started = false;
817 let mut thinking_started = false;
818 let mut tool_indices = std::collections::HashMap::new();
819 let mut reasoning_detail_buffers = std::collections::HashMap::new();
820 let mut inline_reasoning_tags = InlineReasoningTagState::default();
821
822 let outcome = parse_sse_data_frame(
823 " [DONE] ",
824 &mut content_index,
825 &mut text_started,
826 &mut thinking_started,
827 &mut tool_indices,
828 &mut reasoning_detail_buffers,
829 &mut inline_reasoning_tags,
830 ReasoningStreamStyle::SeparateField,
831 );
832
833 assert!(
834 matches!(outcome, SseDataFrame::Done),
835 "`data: [DONE]` must terminate the stream instead of waiting for the HTTP connection to close"
836 );
837 assert_eq!(content_index, 0);
838 assert!(!text_started);
839 assert!(!thinking_started);
840 assert!(tool_indices.is_empty());
841 }
842
843 #[test]
844 fn decoder_emits_tool_use_block_for_tool_call_delta() {
845 // Tool-call deltas are content too — once one arrives, transparent
846 // retry must be off (the model has committed to a tool invocation
847 // path that DeepSeek has billed for).
848 let events = decode_chunk(
849 r#"{"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_1","function":{"name":"grep_files","arguments":"{\"pattern\":\"foo\"}"}}]}}]}"#,
850 );
851 assert!(
852 events.iter().any(|e| matches!(
853 e,
854 StreamEvent::ContentBlockStart {
855 content_block: ContentBlockStart::ToolUse { name, ..},
856 ..
857 } if name == "grep_files"
858 )),
859 "should open a ToolUse block for grep_files; got {events:?}"
860 );
861 assert!(
862 events.iter().any(|e| matches!(
863 e,
864 StreamEvent::ContentBlockDelta {
865 delta: Delta::InputJsonDelta { partial_json },
866 ..
867 } if partial_json.contains("\"pattern\"")
868 )),
869 "should yield InputJsonDelta carrying the tool args; got {events:?}"
870 );
871 }
872
873 #[test]
874 fn decoder_uses_fallback_name_for_empty_streaming_tool_name() {
875 let events = decode_chunk(
876 r#"{"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_empty","function":{"name":"","arguments":"{}"}}]}}]}"#,
877 );
878
879 assert!(
880 events.iter().any(|event| matches!(
881 event,
882 StreamEvent::ContentBlockStart {
883 content_block: ContentBlockStart::ToolUse { name, ..},
884 ..
885 } if name == "unknown_tool"
886 )),
887 "empty upstream tool names should render as unknown_tool; got {events:?}"
888 );
889 }
890
891 #[test]
892 fn non_streaming_response_uses_fallback_name_for_missing_tool_name() {
893 let payload: Value = serde_json::from_str(
894 r#"{
895 "id": "chatcmpl_1",
896 "model": "deepseek-v4-pro",
897 "choices": [{
898 "message": {
899 "role": "assistant",
900 "tool_calls": [{
901 "id": "call_missing",
902 "function": { "arguments": "{}" }
903 }]
904 },
905 "finish_reason": "tool_calls"
906 }]
907 }"#,
908 )
909 .expect("valid response");
910
911 let parsed = parse_chat_message(&payload).expect("message parses");
912 let tool_name = parsed.content.iter().find_map(|block| match block {
913 ContentBlock::ToolUse { name, .. } => Some(name.as_str()),
914 _ => None,
915 });
916
917 assert_eq!(tool_name, Some("unknown_tool"));
918 }
919
920 /// Regression for the parallel-tool-calls-without-id collision (audit
921 /// Finding 8): when the upstream chunk omits the `id` field, the
922 /// fallback used to be the literal string `"tool_call"` for every
923 /// parallel call, so two tool calls in one delta ended up sharing an
924 /// id. Downstream routing then matched the first call's tool_result
925 /// twice and the second call hung. The fallback is now indexed by the
926 /// content-block position, keeping each call unique within the
927 /// response.
928 #[test]
929 fn decoder_assigns_unique_fallback_ids_to_parallel_tool_calls_missing_id() {
930 let events = decode_chunk(
931 r#"{"choices":[{"delta":{"tool_calls":[
932 {"index":0,"function":{"name":"grep_files","arguments":"{\"pattern\":\"a\"}"}},
933 {"index":1,"function":{"name":"read_file","arguments":"{\"path\":\"x\"}"}}
934 ]}}]}"#,
935 );
936
937 let ids: Vec<&str> = events
938 .iter()
939 .filter_map(|e| match e {
940 StreamEvent::ContentBlockStart {
941 content_block: ContentBlockStart::ToolUse { id, .. },
942 ..
943 } => Some(id.as_str()),
944 _ => None,
945 })
946 .collect();
947
948 assert_eq!(
949 ids.len(),
950 2,
951 "expected two tool-use blocks for parallel tool calls; got {events:?}"
952 );
953 assert_ne!(
954 ids[0], ids[1],
955 "parallel tool calls without upstream `id` must get distinct fallback ids; got {ids:?}"
956 );
957 }
958
959 #[test]
960 fn decoder_preserves_upstream_tool_call_id_when_present() {
961 // Counter-test to the fallback regression: when the upstream chunk
962 // does include `id`, we forward it verbatim — we shouldn't quietly
963 // rewrite ids the API gave us just because we have a fallback path.
964 let events = decode_chunk(
965 r#"{"choices":[{"delta":{"tool_calls":[{"index":0,"id":"call_xyz","function":{"name":"grep_files","arguments":"{}"}}]}}]}"#,
966 );
967 let id = events
968 .iter()
969 .find_map(|e| match e {
970 StreamEvent::ContentBlockStart {
971 content_block: ContentBlockStart::ToolUse { id, .. },
972 ..
973 } => Some(id.as_str()),
974 _ => None,
975 })
976 .expect("tool-use block present");
977 assert_eq!(id, "call_xyz");
978 }
979
980 #[test]
981 fn request_builder_preserves_internal_system_messages() {
982 let messages = vec![Message {
983 role: Role::System,
984 content: vec![ContentBlock::Text {
985 text: "internal runtime event".to_string(),
986 cache_control: None,
987 }],
988 }];
989
990 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
991
992 assert_eq!(built.len(), 1);
993 assert_eq!(built[0]["role"], "system");
994 assert_eq!(built[0]["content"], "internal runtime event");
995 }
996
997 fn tool_use_message(id: &str, name: &str, input: Value) -> Message {
998 Message {
999 role: Role::Assistant,
1000 content: vec![ContentBlock::ToolUse {
1001 execution_id: None,
1002 id: id.to_string(),
1003 name: name.to_string(),
1004 input,
1005 caller: None,
1006 thought_signature: None,
1007 }],
1008 }
1009 }
1010
1011 fn tool_result_message(id: &str, content: &str) -> Message {
1012 Message {
1013 role: Role::User,
1014 content: vec![ContentBlock::ToolResult {
1015 execution_id: None,
1016 tool_use_id: id.to_string(),
1017 content: content.to_string(),
1018 is_error: None,
1019 content_blocks: None,
1020 }],
1021 }
1022 }
1023
1024 fn user_message_with_turn_meta(turn_meta: &str, task: &str) -> Message {
1025 Message {
1026 role: Role::User,
1027 content: vec![
1028 ContentBlock::Text {
1029 text: turn_meta.to_string(),
1030 cache_control: None,
1031 },
1032 ContentBlock::Text {
1033 text: task.to_string(),
1034 cache_control: None,
1035 },
1036 ],
1037 }
1038 }
1039
1040 fn user_message_with_tail_turn_meta(task: &str, turn_meta: &str) -> Message {
1041 Message {
1042 role: Role::User,
1043 content: vec![
1044 ContentBlock::Text {
1045 text: task.to_string(),
1046 cache_control: None,
1047 },
1048 ContentBlock::Text {
1049 text: turn_meta.to_string(),
1050 cache_control: None,
1051 },
1052 ],
1053 }
1054 }
1055
1056 fn tool_message_content(messages: &[Value], index: usize) -> &str {
1057 messages
1058 .iter()
1059 .filter(|message| message.get("role").and_then(Value::as_str) == Some("tool"))
1060 .nth(index)
1061 .and_then(|message| message.get("content").and_then(Value::as_str))
1062 .expect("tool message content")
1063 }
1064
1065 fn user_message_content(messages: &[Value], index: usize) -> &str {
1066 messages
1067 .iter()
1068 .filter(|message| message.get("role").and_then(Value::as_str) == Some("user"))
1069 .nth(index)
1070 .and_then(|message| message.get("content").and_then(Value::as_str))
1071 .expect("user message content")
1072 }
1073
1074 fn with_tool_result_sha_spillover_root<T>(f: impl FnOnce() -> T) -> T {
1075 let _guard = crate::tools::truncate::TEST_SPILLOVER_GUARD
1076 .lock()
1077 .unwrap_or_else(|err| err.into_inner());
1078 let tmp = tempfile::tempdir().expect("tempdir");
1079 let prior = crate::tools::truncate::set_test_spillover_root(Some(
1080 tmp.path().join(".deepseek").join("tool_outputs"),
1081 ));
1082 struct Restore(Option<std::path::PathBuf>);
1083 impl Drop for Restore {
1084 fn drop(&mut self) {
1085 crate::tools::truncate::set_test_spillover_root(self.0.take());
1086 }
1087 }
1088 let _restore = Restore(prior);
1089 f()
1090 }
1091
1092 #[test]
1093 fn request_builder_deduplicates_consecutive_identical_turn_meta_for_wire() {
1094 let turn_meta = "<turn_meta>\nCurrent local date: 2026-05-09\n</turn_meta>";
1095 let messages = vec![
1096 user_message_with_turn_meta(turn_meta, "first task"),
1097 Message {
1098 role: Role::Assistant,
1099 content: vec![ContentBlock::Text {
1100 text: "first answer".to_string(),
1101 cache_control: None,
1102 }],
1103 },
1104 user_message_with_turn_meta(turn_meta, "second task"),
1105 ];
1106
1107 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1108 let first = user_message_content(&built, 0);
1109 let second = user_message_content(&built, 1);
1110 let expected_ref = "<turn_meta_unchanged />";
1111
1112 assert!(first.starts_with(turn_meta), "got: {first}");
1113 assert!(second.starts_with(expected_ref), "got: {second}");
1114 assert!(second.ends_with("second task"), "got: {second}");
1115 assert_eq!(
1116 second,
1117 format!("{expected_ref}\nsecond task"),
1118 "ref text must stay stable"
1119 );
1120 }
1121
1122 #[test]
1123 fn request_builder_keeps_tail_turn_meta_after_user_text_for_wire() {
1124 let turn_meta = "<turn_meta>\nCurrent local date: 2026-05-09\n</turn_meta>";
1125 let messages = vec![
1126 user_message_with_tail_turn_meta("first task", turn_meta),
1127 Message {
1128 role: Role::Assistant,
1129 content: vec![ContentBlock::Text {
1130 text: "first answer".to_string(),
1131 cache_control: None,
1132 }],
1133 },
1134 user_message_with_tail_turn_meta("second task", turn_meta),
1135 ];
1136
1137 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1138 let first = user_message_content(&built, 0);
1139 let second = user_message_content(&built, 1);
1140 let expected_ref = "<turn_meta_unchanged />";
1141
1142 assert_eq!(first, format!("first task\n{turn_meta}"));
1143 assert_eq!(second, format!("second task\n{expected_ref}"));
1144 }
1145
1146 #[test]
1147 fn request_builder_keeps_changed_turn_meta_full_and_updates_recent_hash() {
1148 let first_meta = "<turn_meta>\nCurrent local date: 2026-05-09\n</turn_meta>";
1149 let second_meta =
1150 "<turn_meta>\nCurrent local date: 2026-05-09\nWorking set: src/lib.rs\n</turn_meta>";
1151 let messages = vec![
1152 user_message_with_turn_meta(first_meta, "first task"),
1153 user_message_with_turn_meta(second_meta, "second task"),
1154 ];
1155
1156 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1157 let first = user_message_content(&built, 0);
1158 let second = user_message_content(&built, 1);
1159
1160 assert!(first.starts_with(first_meta), "got: {first}");
1161 assert!(second.starts_with(second_meta), "got: {second}");
1162 assert!(!second.contains("<TURN_META_REF"), "got: {second}");
1163 }
1164
1165 #[test]
1166 fn turn_meta_dedup_is_wire_only_and_does_not_mutate_session_message() {
1167 let turn_meta = "<turn_meta>\nCurrent local date: 2026-05-09\n</turn_meta>";
1168 let messages = vec![
1169 user_message_with_turn_meta(turn_meta, "first task"),
1170 user_message_with_turn_meta(turn_meta, "second task"),
1171 ];
1172
1173 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1174 assert!(
1175 user_message_content(&built, 1).starts_with("<turn_meta_unchanged />"),
1176 "got: {}",
1177 user_message_content(&built, 1)
1178 );
1179
1180 match &messages[1].content[0] {
1181 ContentBlock::Text { text, .. } => assert_eq!(text, turn_meta),
1182 other => panic!("expected text block, got {other:?}"),
1183 }
1184 }
1185
1186 #[test]
1187 fn cache_inspect_reports_turn_meta_dedup_metadata() {
1188 let turn_meta = format!(
1189 "<turn_meta>\nCurrent local date: 2026-05-09\n{}\n</turn_meta>",
1190 "Working set: src/lib.rs\n".repeat(20)
1191 );
1192 let request = MessageRequest {
1193 model: "deepseek-v4-flash".to_string(),
1194 messages: vec![
1195 user_message_with_turn_meta(&turn_meta, "first task"),
1196 user_message_with_turn_meta(&turn_meta, "second task"),
1197 ],
1198 max_tokens: 0,
1199 system: None,
1200 tools: None,
1201 tool_choice: None,
1202 metadata: None,
1203 thinking: None,
1204 reasoning_effort: None,
1205 stream: None,
1206 temperature: None,
1207 top_p: None,
1208 };
1209
1210 let inspection = inspect_prompt_for_request(&request);
1211 let turn_meta_layers: Vec<_> = inspection
1212 .layers
1213 .iter()
1214 .filter_map(|layer| layer.turn_meta.as_ref())
1215 .collect();
1216
1217 assert_eq!(turn_meta_layers.len(), 2);
1218 assert_eq!(
1219 turn_meta_layers[0].original_chars,
1220 turn_meta.chars().count()
1221 );
1222 assert_eq!(turn_meta_layers[0].sent_chars, turn_meta.chars().count());
1223 assert!(!turn_meta_layers[0].deduplicated);
1224 assert_eq!(turn_meta_layers[0].sha256, sha256_hex(turn_meta.as_bytes()));
1225 assert_eq!(
1226 turn_meta_layers[1].original_chars,
1227 turn_meta.chars().count()
1228 );
1229 assert!(turn_meta_layers[1].sent_chars < turn_meta_layers[1].original_chars);
1230 assert!(turn_meta_layers[1].deduplicated);
1231 assert_eq!(turn_meta_layers[1].sha256, turn_meta_layers[0].sha256);
1232 }
1233
1234 #[test]
1235 fn request_builder_truncates_large_tool_result_for_wire() {
1236 let long_output = format!("{}{}", "A".repeat(70_000), "Z".repeat(70_000));
1237 let messages = vec![
1238 tool_use_message(
1239 "tool-long",
1240 "shell_command",
1241 json!({"command": "cargo test"}),
1242 ),
1243 tool_result_message("tool-long", &long_output),
1244 ];
1245
1246 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1247 let sent = tool_message_content(&built, 0);
1248
1249 assert!(sent.contains("[TOOL_RESULT_TRUNCATED]"), "got: {sent}");
1250 assert!(sent.contains("tool_name: shell_command"), "got: {sent}");
1251 assert!(sent.contains("command_or_query: cargo test"), "got: {sent}");
1252 assert!(sent.contains("original_chars: 140000"), "got: {sent}");
1253 assert!(sent.contains("sha256:"), "got: {sent}");
1254 assert!(
1255 sent.contains("exact_detail: unavailable; no session-owned artifact was recorded"),
1256 "got: {sent}"
1257 );
1258 assert!(!sent.contains("retrieve_tool_result"), "got: {sent}");
1259 // The model catalog supplies the window when no explicit route is available.
1260 let budget = tool_result_sent_char_budget("deepseek-v4-flash");
1261 assert!(budget >= 100_000);
1262 let excerpt = budget - TOOL_RESULT_EXCERPT_FRAME_CHARS;
1263 let head = excerpt * 2 / 3;
1264 assert!(sent.contains(&"A".repeat(head)), "head kept");
1265 assert!(sent.contains(&"Z".repeat(excerpt - head)), "tail kept");
1266 assert!(
1267 sent.contains(&format!(
1268 "truncated {} chars from middle",
1269 140_000 - excerpt
1270 )),
1271 "omitted count"
1272 );
1273 assert_ne!(sent, long_output);
1274 }
1275
1276 #[test]
1277 fn restored_tool_results_obey_the_active_route_budget() {
1278 let output = "x".repeat(20_000);
1279 let recovery = format!(
1280 "{output}\n{}",
1281 crate::tools::truncate::SPILLOVER_RECOVERY_HINT
1282 );
1283 let request = MessageRequest {
1284 model: "trinity-mini".into(),
1285 messages: vec![
1286 tool_use_message("raw", "read_file", json!({"path": "a.rs"})),
1287 tool_result_message("raw", &output),
1288 tool_use_message("repeat", "read_file", json!({"path": "a.rs"})),
1289 tool_result_message("repeat", &output),
1290 tool_use_message("saved", "read_file", json!({"path": "b.rs"})),
1291 tool_result_message("saved", &recovery),
1292 ],
1293 max_tokens: 1_024,
1294 system: None,
1295 tools: None,
1296 tool_choice: None,
1297 metadata: None,
1298 thinking: None,
1299 reasoning_effort: None,
1300 stream: None,
1301 temperature: None,
1302 top_p: None,
1303 };
1304 // The bundled 128K route and an explicit smaller offering both beat
1305 // the old unknown-route 100K ceiling. Streaming and blocking share it.
1306 for (limits, expected_budget) in [
1307 (None, 15_360),
1308 (
1309 Some(RouteLimits {
1310 context_tokens: Some(32_000),
1311 ..RouteLimits::default()
1312 }),
1313 3_840,
1314 ),
1315 ] {
1316 for stream in [false, true] {
1317 let wire = build_chat_wire_body(
1318 &request,
1319 ProviderKind::Arcee,
1320 "https://api.arcee.ai/api/v1",
1321 stream,
1322 limits,
1323 )
1324 .unwrap();
1325 let messages = wire.body["messages"].as_array().unwrap();
1326 for index in [0, 1] {
1327 let sent = tool_message_content(messages, index);
1328 assert!(sent.contains("[TOOL_RESULT_TRUNCATED]"));
1329 assert!(sent.chars().count() <= expected_budget);
1330 assert!(!sent.contains("<TOOL_RESULT_REF"));
1331 assert!(sent.contains("exact_detail: unavailable"));
1332 }
1333 assert_eq!(tool_message_content(messages, 2), recovery);
1334 }
1335 }
1336 assert!(matches!(&request.messages[1].content[0],
1337 ContentBlock::ToolResult { content, .. } if content == &output));
1338 }
1339
1340 #[test]
1341 fn request_builder_keeps_unowned_extreme_tool_output_bounded_without_false_hint() {
1342 with_tool_result_sha_spillover_root(|| {
1343 let huge_output = format!(
1344 "{}{}{}",
1345 "DIFF_HEAD\n".repeat(10_000),
1346 "MIDDLE_POISON\n".repeat(10_000),
1347 "DIFF_TAIL\n".repeat(10_000)
1348 );
1349 let sha = sha256_hex(huge_output.as_bytes());
1350 let messages = vec![
1351 tool_use_message("tool-huge", "exec_shell", json!({"command": "git diff"})),
1352 tool_result_message("tool-huge", &huge_output),
1353 ];
1354
1355 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1356 let sent = tool_message_content(&built, 0);
1357
1358 assert!(sent.contains("[TOOL_RESULT_TRUNCATED]"), "got: {sent}");
1359 assert!(sent.contains("tool_name: exec_shell"), "got: {sent}");
1360 assert!(sent.contains("command_or_query: git diff"), "got: {sent}");
1361 assert!(sent.contains(&format!("sha256: {sha}")), "got: {sent}");
1362 assert!(sent.contains("exact_detail: unavailable"), "got: {sent}");
1363 assert!(!sent.contains("retrieve_tool_result"), "got: {sent}");
1364 assert!(
1365 sent.chars().count() <= tool_result_sent_char_budget("deepseek-v4-flash"),
1366 "truncated result should stay bounded, sent {} chars",
1367 sent.chars().count()
1368 );
1369 assert!(
1370 !sent.contains("MIDDLE_POISON"),
1371 "omitted middle should not be sent to the next model turn"
1372 );
1373 assert_ne!(sent, huge_output);
1374 });
1375 }
1376
1377 #[test]
1378 fn request_builder_does_not_dedup_short_tool_results_for_wire() {
1379 let output = "same tool output";
1380 let messages = vec![
1381 tool_use_message("tool-1", "read_file", json!({"path": "README.md"})),
1382 tool_result_message("tool-1", output),
1383 tool_use_message("tool-2", "read_file", json!({"path": "README.md"})),
1384 tool_result_message("tool-2", output),
1385 ];
1386
1387 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1388 let first = tool_message_content(&built, 0);
1389 let second = tool_message_content(&built, 1);
1390
1391 assert_eq!(first, output);
1392 assert_eq!(second, output);
1393 assert!(!second.contains("<TOOL_RESULT_REF"), "got: {second}");
1394 }
1395
1396 #[test]
1397 fn request_builder_deduplicates_medium_identical_tool_results_to_earlier_message() {
1398 with_tool_result_sha_spillover_root(|| {
1399 // 2,000 chars is intentionally above TOOL_RESULT_DEDUP_MIN_CHARS
1400 // (1,024) but below the wire backstop budget. This
1401 // verifies the cache-saving path for repeated medium outputs that
1402 // do not otherwise need truncation.
1403 let output = "A".repeat(2_000);
1404 let messages = vec![
1405 tool_use_message("tool-1", "read_file", json!({"path": "README.md"})),
1406 tool_result_message("tool-1", &output),
1407 tool_use_message("tool-2", "read_file", json!({"path": "README.md"})),
1408 tool_result_message("tool-2", &output),
1409 ];
1410
1411 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1412 let first = tool_message_content(&built, 0);
1413 let second = tool_message_content(&built, 1);
1414
1415 assert_eq!(first, output);
1416 assert!(!first.contains("[TOOL_RESULT_TRUNCATED]"), "got: {first}");
1417 assert!(
1418 second.starts_with("<TOOL_RESULT_REF sha=\""),
1419 "got: {second}"
1420 );
1421 assert!(
1422 second.contains("original_message=\"Message #1\""),
1423 "got: {second}"
1424 );
1425 assert!(second.contains("chars=\"2000\""), "got: {second}");
1426 assert!(
1427 second.contains("source: full content appears in Message #1 earlier in this request"),
1428 "got: {second}"
1429 );
1430 assert!(!second.contains("retrieve_tool_result"), "got: {second}");
1431 });
1432 }
1433
1434 #[test]
1435 fn request_builder_never_dedups_large_identical_write_file_confirmations() {
1436 with_tool_result_sha_spillover_root(|| {
1437 // A `write_file` result embeds the unified diff + summary; it is a
1438 // confirmation, not retrievable data. Two identical >1024-char
1439 // write_file results must BOTH stay inline — collapsing the second
1440 // to a SHA ref makes the model lose write-success context and
1441 // report the file as missing (#1695).
1442 let output = "A".repeat(2_000);
1443 let messages = vec![
1444 tool_use_message("tool-1", "write_file", json!({"path": "big.txt"})),
1445 tool_result_message("tool-1", &output),
1446 tool_use_message("tool-2", "write_file", json!({"path": "big.txt"})),
1447 tool_result_message("tool-2", &output),
1448 ];
1449
1450 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1451 let first = tool_message_content(&built, 0);
1452 let second = tool_message_content(&built, 1);
1453
1454 assert_eq!(first, output);
1455 assert_eq!(second, output);
1456 assert!(!second.contains("<TOOL_RESULT_REF"), "got: {second}");
1457
1458 // Non-mutation tools still dedup: an identical medium read_file
1459 // result points back to the first full message in this request.
1460 let read_messages = vec![
1461 tool_use_message("read-1", "read_file", json!({"path": "README.md"})),
1462 tool_result_message("read-1", &output),
1463 tool_use_message("read-2", "read_file", json!({"path": "README.md"})),
1464 tool_result_message("read-2", &output),
1465 ];
1466 let read_built = build_chat_messages(None, &read_messages, "deepseek-v4-flash");
1467 let read_first = tool_message_content(&read_built, 0);
1468 let read_second = tool_message_content(&read_built, 1);
1469 assert_eq!(read_first, output);
1470 assert!(
1471 read_second.starts_with("<TOOL_RESULT_REF sha=\""),
1472 "got: {read_second}"
1473 );
1474 assert!(read_second.contains("source: full content appears in Message #1"));
1475 assert!(!read_second.contains("retrieve_tool_result"));
1476 });
1477 }
1478
1479 #[test]
1480 fn large_unowned_results_stay_bounded_without_false_retrieval_handles() {
1481 // The adaptive router normally replaces a large result with a
1482 // session-owned artifact receipt before this provider-wire fallback.
1483 // If legacy/raw history reaches here, it may be excerpted but must not
1484 // advertise the process-wide SHA store as retrievable.
1485 let big_diff = "D".repeat(120_000);
1486 let sha = sha256_hex(big_diff.as_bytes());
1487
1488 let messages = vec![
1489 tool_use_message("w-1", "write_file", json!({"path": "huge.rs"})),
1490 tool_result_message("w-1", &big_diff),
1491 tool_use_message("w-2", "write_file", json!({"path": "huge.rs"})),
1492 tool_result_message("w-2", &big_diff),
1493 ];
1494 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1495 let first = tool_message_content(&built, 0);
1496 let second = tool_message_content(&built, 1);
1497
1498 // Mutation confirmations are independently excerpted, never deduped.
1499 assert!(
1500 first.contains("[TOOL_RESULT_TRUNCATED]"),
1501 "first should be truncated, got: {first}"
1502 );
1503 assert!(
1504 !first.contains("<TOOL_RESULT_REF"),
1505 "first must not be a dedup ref, got: {first}"
1506 );
1507 assert!(
1508 !second.contains("<TOOL_RESULT_REF"),
1509 "second identical write_file must stay inline (#1695), got: {second}"
1510 );
1511 assert!(
1512 second.contains("[TOOL_RESULT_TRUNCATED]"),
1513 "second should also be inline-truncated, got: {second}"
1514 );
1515 assert!(
1516 first.contains(&format!("sha256: {sha}")),
1517 "truncation block should retain an integrity digest, got: {first}"
1518 );
1519 assert!(first.contains("exact_detail: unavailable"));
1520 assert!(!first.contains("retrieve_tool_result"));
1521
1522 // A huge non-mutation result cannot refer to an earlier *full* message,
1523 // because both wire messages are excerpts. It therefore stays a
1524 // truthful bounded excerpt too.
1525 let read_messages = vec![
1526 tool_use_message("r-1", "read_file", json!({"path": "huge.rs"})),
1527 tool_result_message("r-1", &big_diff),
1528 tool_use_message("r-2", "read_file", json!({"path": "huge.rs"})),
1529 tool_result_message("r-2", &big_diff),
1530 ];
1531 let read_built = build_chat_messages(None, &read_messages, "deepseek-v4-flash");
1532 let read_second = tool_message_content(&read_built, 1);
1533 assert!(read_second.contains("[TOOL_RESULT_TRUNCATED]"));
1534 assert!(!read_second.contains("<TOOL_RESULT_REF"));
1535 assert!(!read_second.contains("retrieve_tool_result"));
1536 }
1537
1538 #[test]
1539 fn tool_result_budget_is_wire_only_and_does_not_mutate_session_message() {
1540 let long_output = format!("{}{}", "A".repeat(70_000), "Z".repeat(70_000));
1541 let messages = vec![
1542 tool_use_message(
1543 "tool-long",
1544 "shell_command",
1545 json!({"command": "cargo test"}),
1546 ),
1547 tool_result_message("tool-long", &long_output),
1548 ];
1549
1550 let built = build_chat_messages(None, &messages, "deepseek-v4-flash");
1551 let sent = tool_message_content(&built, 0);
1552 assert_ne!(sent, long_output);
1553
1554 match &messages[1].content[0] {
1555 ContentBlock::ToolResult { content, .. } => assert_eq!(content, &long_output),
1556 other => panic!("expected tool result, got {other:?}"),
1557 }
1558 }
1559
1560 #[test]
1561 fn cache_inspect_reports_bounded_unowned_tool_result_metadata() {
1562 let long_output = format!("{}{}", "A".repeat(70_000), "Z".repeat(70_000));
1563 let request = MessageRequest {
1564 model: "deepseek-v4-flash".to_string(),
1565 messages: vec![
1566 tool_use_message("tool-1", "shell_command", json!({"command": "cargo test"})),
1567 tool_result_message("tool-1", &long_output),
1568 tool_use_message("tool-2", "shell_command", json!({"command": "cargo test"})),
1569 tool_result_message("tool-2", &long_output),
1570 ],
1571 max_tokens: 0,
1572 system: None,
1573 tools: None,
1574 tool_choice: None,
1575 metadata: None,
1576 thinking: None,
1577 reasoning_effort: None,
1578 stream: None,
1579 temperature: None,
1580 top_p: None,
1581 };
1582
1583 let inspection = inspect_prompt_for_request(&request);
1584 let tool_layers: Vec<_> = inspection
1585 .layers
1586 .iter()
1587 .filter_map(|layer| layer.tool_result.as_ref())
1588 .collect();
1589
1590 assert_eq!(tool_layers.len(), 2);
1591 for layer in tool_layers {
1592 assert_eq!(layer.original_chars, 140_000);
1593 assert!(layer.sent_chars < layer.original_chars);
1594 assert!(layer.truncated);
1595 assert!(!layer.deduplicated);
1596 }
1597 }
1598
1599 #[test]
1600 fn mistral_stream_blocks_are_decoded_only_by_the_mistral_style() {
1601 let chunk = r#"{
1602 "choices": [{
1603 "index": 0,
1604 "delta": {"content": [
1605 {"type": "thinking", "thinking": [
1606 {"type": "text", "text": "private trace"}
1607 ], "closed": true},
1608 {"type": "text", "text": "public answer"}
1609 ]},
1610 "finish_reason": null
1611 }]
1612 }"#;
1613
1614 let mistral = decode_chunks_with_style(&[chunk], ReasoningStreamStyle::MistralBlocks);
1615 assert!(mistral.iter().any(|event| matches!(
1616 event,
1617 StreamEvent::ContentBlockDelta {
1618 delta: Delta::ThinkingDelta { thinking },
1619 ..
1620 } if thinking == "private trace"
1621 )));
1622 assert!(mistral.iter().any(|event| matches!(
1623 event,
1624 StreamEvent::ContentBlockDelta {
1625 delta: Delta::TextDelta { text },
1626 ..
1627 } if text == "public answer"
1628 )));
1629
1630 let generic = decode_chunks_with_style(&[chunk], ReasoningStreamStyle::None);
1631 assert!(!generic.iter().any(|event| matches!(
1632 event,
1633 StreamEvent::ContentBlockDelta {
1634 delta: Delta::ThinkingDelta { .. },
1635 ..
1636 }
1637 )));
1638 assert!(!generic.iter().any(|event| matches!(
1639 event,
1640 StreamEvent::ContentBlockDelta {
1641 delta: Delta::TextDelta { .. },
1642 ..
1643 }
1644 )));
1645 }
1646
1647 #[test]
1648 fn deepseek_flash_v41_classifies_reasoning_through_the_catalog() {
1649 // #6044: V4.1's official id dropped the version number, so the literal
1650 // `deepseek-v4` arms cannot see it. The catalog owns the capability and
1651 // every classifier — stream style, wire replay, prompt inspection — must
1652 // read it there instead of relying on another classifier's fallback.
1653 let base_url = "https://api.deepseek.com";
1654 assert!(
1655 requires_reasoning_content("deepseek-flash"),
1656 "the name gate must recognize the official V4.1 id through the catalog"
1657 );
1658 assert!(
1659 should_replay_reasoning_content("deepseek-flash", None),
1660 "prompt inspection must agree with the wire request"
1661 );
1662 assert!(should_replay_reasoning_content_for_provider_on_route(
1663 ProviderKind::Deepseek,
1664 base_url,
1665 "deepseek-flash",
1666 None,
1667 ));
1668
1669 let style =
1670 reasoning_stream_style_for_route(ProviderKind::Deepseek, base_url, "deepseek-flash", None);
1671 assert_eq!(style, ReasoningStreamStyle::SeparateField);
1672 let events = decode_chunks_with_style(
1673 &[r#"{"choices":[{"delta":{"reasoning_content":"private flash plan"}}]}"#],
1674 style,
1675 );
1676 assert_eq!(thinking_delta_text(&events), "private flash plan");
1677 assert_eq!(
1678 text_delta_text(&events),
1679 "",
1680 "reasoning must never leak into visible prose"
1681 );
1682
1683 // A non-reasoning DeepSeek name stays literal: the prefix alone is not
1684 // evidence, the catalog entry is.
1685 assert!(!requires_reasoning_content("deepseek-coder"));
1686 }
1687
1687 lines RUST