From 464489990a488dd3dbc15470c177e1393d7b4411 Mon Sep 17 00:00:00 2001 From: yechank <161688079+yechank-nvidia@users.noreply.github.com> Date: Wed, 30 Sep 2026 05:37:29 +0000 Subject: [PATCH 1/5] fix(grpc): stop the open Messages content block before another starts Messages streaming stopped an open text or tool_use block before a new tool_use block, and stopped an open thinking block only when reasoning ended in a chunk that also carried text. A new thinking or text block stopped nothing. The new block then started at the index of the block that was still open, and the deltas and stops that followed went to an index that was never started. With the existing parsers this happens when a specific tool_choice follows reasoning whose last chunk has no text after `` (deepseek_r1 with the deepseek tool parser, or qwen3 with qwen when the output opens with ``), and when a reasoning parser enters reasoning again after text or a tool call (deepseek_v41, inkling), also within one chunk. Stop whichever block is open, and move to the next index, before every content_block_start. Signed-off-by: yechank <161688079+yechank-nvidia@users.noreply.github.com> --- .../src/routers/grpc/regular/streaming.rs | 158 +++++++---- .../grpc/regular/streaming/eof_tests.rs | 260 ++++++++++++++++++ 2 files changed, 358 insertions(+), 60 deletions(-) diff --git a/model_gateway/src/routers/grpc/regular/streaming.rs b/model_gateway/src/routers/grpc/regular/streaming.rs index fa22f92c1..8d1fdc7f0 100644 --- a/model_gateway/src/routers/grpc/regular/streaming.rs +++ b/model_gateway/src/routers/grpc/regular/streaming.rs @@ -1823,6 +1823,25 @@ impl StreamingProcessor { .map_err(|_| "Client disconnected".to_string()) } + /// Stop the open content block, if any, so the next block gets the next + /// index: reasoning, text and tool calls can alternate. + async fn stop_open_block( + tx: &SseSender, + buffer: &mut Vec, + index: &mut u32, + open: [&mut bool; 3], + ) -> Result<(), String> { + if open + .into_iter() + .fold(false, |any, open| std::mem::take(open) | any) + { + let stop = MessageStreamEvent::ContentBlockStop { index: *index }; + Self::send_messages_event(tx, buffer, &stop).await?; + *index += 1; + } + Ok(()) + } + /// Process reasoning content in Messages streaming mode (n=1 only). /// /// Returns `(normal_text, reasoning_text, in_reasoning)`. @@ -2275,6 +2294,17 @@ impl StreamingProcessor { // Emit thinking content block deltas if !reasoning_chunk_text.is_empty() { if !thinking_block_open { + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; Self::send_messages_event( tx, &mut sse_buffer, @@ -2326,19 +2356,18 @@ impl StreamingProcessor { // Specific function: entire output is arguments for one tool if !has_tool_calls { has_tool_calls = true; - // Close text block if open before starting tool block - if text_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockStop { - index: current_block_index, - }, - ) - .await?; - text_block_open = false; - current_block_index += 1; - } + // Close the open block before starting tool block + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; // Emit content_block_start for the tool_use let tool_name = match &original_request.tool_choice { Some(messages::ToolChoice::Tool { name, .. }) => name.clone(), @@ -2389,6 +2418,17 @@ impl StreamingProcessor { // Emit normal text from parser as text content blocks if !text.is_empty() { if !text_block_open { + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; Self::send_messages_event( tx, &mut sse_buffer, @@ -2420,29 +2460,17 @@ impl StreamingProcessor { if let Some(ref name) = tool_call_item.name { // New tool call: close previous blocks, emit start - if text_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockStop { - index: current_block_index, - }, - ) - .await?; - text_block_open = false; - current_block_index += 1; - } - if tool_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockStop { - index: current_block_index, - }, - ) - .await?; - current_block_index += 1; - } + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; let tool_call_id = utils::generate_tool_call_id( model, @@ -2495,6 +2523,17 @@ impl StreamingProcessor { // Regular text emission (no tools active) if !normal_text.is_empty() { if !text_block_open { + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; Self::send_messages_event( tx, &mut sse_buffer, @@ -2529,6 +2568,17 @@ impl StreamingProcessor { let leftover_text = parser.take_unstreamed_normal_text(); if !leftover_text.is_empty() { if !text_block_open { + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; Self::send_messages_event( tx, &mut sse_buffer, @@ -2563,30 +2613,18 @@ impl StreamingProcessor { has_tool_calls = true; if let Some(ref name) = tool_call_item.name { - // Close text block if open before starting tool block - if text_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockStop { - index: current_block_index, - }, - ) - .await?; - text_block_open = false; - current_block_index += 1; - } - if tool_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockStop { - index: current_block_index, - }, - ) - .await?; - current_block_index += 1; - } + // Close the open block before starting tool block + Self::stop_open_block( + tx, + &mut sse_buffer, + &mut current_block_index, + [ + &mut thinking_block_open, + &mut text_block_open, + &mut tool_block_open, + ], + ) + .await?; let tool_call_id = utils::generate_tool_call_id( model, diff --git a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs index 2787a00d0..4b6eb72f2 100644 --- a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs +++ b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs @@ -15,6 +15,7 @@ use smg_grpc_client::vllm_engine::{ }; use tokio::{net::TcpListener, task::JoinHandle}; use tonic::codec::Codec; +use tool_parser::types::ToolCallItem; use super::*; use crate::{routers::common::sse::SseReceiver, worker::WorkerRegistry}; @@ -578,3 +579,262 @@ async fn deepseek_does_not_emit_aggregate_usage_without_complete_frames() { ); } } + +/// Reasoning for every chunk but "text"; `is_in_reasoning` reports `.0`. +struct ReasoningButText(bool); + +impl ReasoningParser for ReasoningButText { + fn detect_and_parse_reasoning( + &mut self, + text: &str, + ) -> Result { + Ok(ParserResult::normal(text.to_string())) + } + + fn parse_reasoning_streaming_incremental( + &mut self, + text: &str, + ) -> Result { + Ok(match text { + "text" => ParserResult::normal(text.to_string()), + _ => ParserResult::reasoning(text.to_string()), + }) + } + + fn reset(&mut self) {} + + fn model_type(&self) -> &str { + "reasoning-but-text" + } + + fn is_in_reasoning(&self) -> bool { + self.0 + } + + fn mark_reasoning_started(&mut self) {} + + fn mark_think_start_stripped(&mut self) {} +} + +/// Reports one whole call on its first chunk, before the chunk's text. +#[derive(Default)] +struct CallFirst { + called: bool, +} + +#[async_trait::async_trait] +impl ToolParser for CallFirst { + async fn parse_complete( + &self, + output: &str, + ) -> tool_parser::errors::ParserResult<(String, Vec)> { + Ok((output.to_string(), Vec::new())) + } + + async fn parse_incremental( + &mut self, + chunk: &str, + _tools: &[Tool], + ) -> tool_parser::errors::ParserResult { + let calls = (!std::mem::replace(&mut self.called, true)).then(|| ToolCallItem { + tool_index: 0, + name: Some("lookup".to_string()), + parameters: "{}".to_string(), + }); + Ok(StreamingParseResult { + normal_text: chunk.to_string(), + calls: calls.into_iter().collect(), + }) + } + + fn has_tool_markers(&self, _text: &str) -> bool { + false + } +} + +/// A processor with the stub parsers above. +fn stub_processor(in_reasoning: bool) -> StreamingProcessor { + let reasoning = ReasoningParserFactory::new(); + reasoning + .registry() + .register_parser("reasoning-but-text", move || { + Box::new(ReasoningButText(in_reasoning)) + }); + let tools = ToolParserFactory::new(); + tools + .registry() + .register_parser("call-first", || Box::new(CallFirst::default())); + let resolver = utils::ParserResolver::new( + Arc::new(WorkerRegistry::new()), + Some("call-first".to_string()), + Some("reasoning-but-text".to_string()), + ); + StreamingProcessor::new(tools, reasoning, resolver, "vllm") +} + +/// The content block starts and stops of a Messages stream of `texts`. A +/// block starts only when none is open, and every delta goes to the open one. +async fn messages_blocks( + processor: StreamingProcessor, + tool_choice: Option, + texts: &[&str], +) -> Vec { + let spec = MessagesResponseSpec { + thinking: Some(messages::ThinkingConfig::Enabled { + budget_tokens: 1024, + display: None, + }), + tool_choice, + has_tools: true, + history_tool_calls_count: 0, + chat_tools: chat_spec(true).tools.unwrap(), + stop_sequences: None, + }; + let mut frames: Vec<_> = texts.iter().map(|text| chunk(0, text)).collect(); + frames.push(complete(0, "stop")); + let (stream, server) = scripted_stream(frames, "0").await; + let (tx, rx) = sse_channel(); + let result = processor + .process_messages_streaming_chunks( + stream, + dispatch(), + Arc::new(CharacterTokenizer::default()), + (None, None, false, false, false), + spec, + &tx, + None, + ) + .await; + drop(tx); + let events = collect_events(rx).await; + server.abort(); + assert!(result.is_ok(), "{result:?}"); + let mut open = None; + let mut blocks = Vec::new(); + for event in &events { + let index = &event["index"]; + match event["type"].as_str() { + Some("content_block_start") => { + blocks.push(format!("start {index} {}", event["content_block"]["type"])); + assert!(open.replace(index).is_none(), "{blocks:?}"); + } + Some("content_block_delta") => { + assert_eq!(open, Some(index), "delta {index} after {blocks:?}"); + } + Some("content_block_stop") => { + blocks.push(format!("stop {index}")); + assert_eq!(open.take(), Some(index), "{blocks:?}"); + } + _ => {} + } + } + blocks +} + +#[tokio::test] +async fn messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternate() { + // Reasoning, then a call with no text between; text; reasoning again. + assert_eq!( + messages_blocks(stub_processor(false), None, &["a", "text", "b"]).await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"tool_use\"", + "stop 1", + "start 2 \"text\"", + "stop 2", + "start 3 \"thinking\"", + "stop 3", + ] + ); + // Text the reasoning parser returns while it stays in reasoning, as when + // one chunk ends reasoning, adds text and starts reasoning again. + assert_eq!( + messages_blocks(stub_processor(true), None, &["a", "text", "b"]).await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"text\"", + "stop 1", + "start 2 \"thinking\"", + "stop 2", + ] + ); + // A specific tool right after reasoning. + let tool = messages::ToolChoice::Tool { + name: "lookup".to_string(), + disable_parallel_tool_use: None, + }; + assert_eq!( + messages_blocks(stub_processor(false), Some(tool), &["a"]).await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"tool_use\"", + "stop 1", + ] + ); +} + +#[tokio::test] +async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { + let processor = |tool: &str, reasoning: &str| { + StreamingProcessor::new( + ToolParserFactory::new(), + ReasoningParserFactory::new(), + utils::ParserResolver::new( + Arc::new(WorkerRegistry::new()), + Some(tool.to_string()), + Some(reasoning.to_string()), + ), + "vllm", + ) + }; + // `` ends reasoning without text, and the arguments of a specific + // tool follow in the next chunk. + let tool = messages::ToolChoice::Tool { + name: "lookup".to_string(), + disable_parallel_tool_use: None, + }; + assert_eq!( + messages_blocks( + processor("deepseek", "deepseek_r1"), + Some(tool), + &["plan", "", "{}"], + ) + .await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"tool_use\"", + "stop 1", + ] + ); + // A reasoning parser that enters reasoning again after text. + assert_eq!( + messages_blocks( + processor("deepseek_v41", "deepseek_v41"), + None, + &[ + "plan", + "", + "answer", + "", + "more", + "", + "done" + ], + ) + .await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"text\"", + "stop 1", + "start 2 \"thinking\"", + "stop 2", + "start 3 \"text\"", + "stop 3", + ] + ); +} From 0c2d362ff2469298d367823bd43a3c495c6ae703 Mon Sep 17 00:00:00 2001 From: yechank <161688079+yechank-nvidia@users.noreply.github.com> Date: Wed, 30 Sep 2026 08:29:45 +0000 Subject: [PATCH 2/5] fix(grpc): keep open tool call arguments ahead of text in Messages streaming With multi-token chunks, a tool parser can return the arguments that finish the open call together with the text after the call. qwen_xml does this for a chunk like `1\n\n\n\n`. Messages streaming emitted the text first. Since the previous commit, starting that text block stops the open tool_use block, so the arguments went to the text block and clients built the tool input without them. Send the calls before the first one with a name, which continue the open call, to the tool_use block before the text. The test helper now also checks that each delta matches the type of its block, and returns each tool input as a client SDK builds it. Signed-off-by: yechank <161688079+yechank-nvidia@users.noreply.github.com> --- .../src/routers/grpc/regular/streaming.rs | 28 +++- .../grpc/regular/streaming/eof_tests.rs | 139 +++++++++++++++--- 2 files changed, 145 insertions(+), 22 deletions(-) diff --git a/model_gateway/src/routers/grpc/regular/streaming.rs b/model_gateway/src/routers/grpc/regular/streaming.rs index 8d1fdc7f0..7c3e98e5e 100644 --- a/model_gateway/src/routers/grpc/regular/streaming.rs +++ b/model_gateway/src/routers/grpc/regular/streaming.rs @@ -2413,8 +2413,34 @@ impl StreamingProcessor { match parser.parse_incremental(&normal_text, chat_tools).await { Ok(StreamingParseResult { normal_text: text, - calls, + mut calls, }) => { + // Arguments that finish the open call come before + // the text after it in the same chunk. + let finishing = if tool_block_open { + calls + .iter() + .position(|call| call.name.is_some()) + .unwrap_or(calls.len()) + } else { + 0 + }; + for tool_call_item in calls.drain(..finishing) { + if !tool_call_item.parameters.is_empty() { + Self::send_messages_event( + tx, + &mut sse_buffer, + &MessageStreamEvent::ContentBlockDelta { + index: current_block_index, + delta: ContentBlockDelta::InputJsonDelta { + partial_json: tool_call_item.parameters, + }, + }, + ) + .await?; + } + } + // Emit normal text from parser as text content blocks if !text.is_empty() { if !text_block_open { diff --git a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs index 4b6eb72f2..945a6cc80 100644 --- a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs +++ b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs @@ -672,13 +672,26 @@ fn stub_processor(in_reasoning: bool) -> StreamingProcessor { StreamingProcessor::new(tools, reasoning, resolver, "vllm") } -/// The content block starts and stops of a Messages stream of `texts`. A -/// block starts only when none is open, and every delta goes to the open one. +/// The content block starts and stops from `messages_blocks_and_inputs`. async fn messages_blocks( processor: StreamingProcessor, tool_choice: Option, texts: &[&str], ) -> Vec { + messages_blocks_and_inputs(processor, tool_choice, texts) + .await + .0 +} + +/// The content block starts and stops of a Messages stream of `texts`, and +/// the input of each `tool_use` block joined from its `input_json_delta`s as +/// a client SDK builds it. A block starts only when none is open, and every +/// delta goes to the open block and matches its type. +async fn messages_blocks_and_inputs( + processor: StreamingProcessor, + tool_choice: Option, + texts: &[&str], +) -> (Vec, Vec) { let spec = MessagesResponseSpec { thinking: Some(messages::ThinkingConfig::Enabled { budget_tokens: 1024, @@ -711,24 +724,46 @@ async fn messages_blocks( assert!(result.is_ok(), "{result:?}"); let mut open = None; let mut blocks = Vec::new(); + let mut inputs = Vec::new(); for event in &events { let index = &event["index"]; match event["type"].as_str() { Some("content_block_start") => { - blocks.push(format!("start {index} {}", event["content_block"]["type"])); - assert!(open.replace(index).is_none(), "{blocks:?}"); + let kind = event["content_block"]["type"].as_str().unwrap_or_default(); + blocks.push(format!("start {index} {kind:?}")); + assert!(open.replace((index, kind)).is_none(), "{blocks:?}"); + if kind == "tool_use" { + inputs.push(String::new()); + } } Some("content_block_delta") => { - assert_eq!(open, Some(index), "delta {index} after {blocks:?}"); + let delta = &event["delta"]; + let kind = match delta["type"].as_str() { + Some("text_delta") => "text", + Some("input_json_delta") => "tool_use", + Some("thinking_delta" | "signature_delta") => "thinking", + other => panic!("unexpected delta {other:?}"), + }; + assert_eq!(open, Some((index, kind)), "delta {delta} after {blocks:?}"); + if let Some(json) = delta["partial_json"].as_str() { + inputs.last_mut().expect("open tool_use").push_str(json); + } } Some("content_block_stop") => { blocks.push(format!("stop {index}")); - assert_eq!(open.take(), Some(index), "{blocks:?}"); + assert_eq!(open.take().map(|open| open.0), Some(index), "{blocks:?}"); } _ => {} } } - blocks + let inputs = inputs + .iter() + .map(|json| match json.as_str() { + "" => serde_json::json!({}), + json => serde_json::from_str(json).expect("tool input JSON"), + }) + .collect(); + (blocks, inputs) } #[tokio::test] @@ -776,20 +811,22 @@ async fn messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternate( ); } +/// A processor with the registered parsers `tool` and `reasoning`. +fn named_processor(tool: &str, reasoning: &str) -> StreamingProcessor { + StreamingProcessor::new( + ToolParserFactory::new(), + ReasoningParserFactory::new(), + utils::ParserResolver::new( + Arc::new(WorkerRegistry::new()), + Some(tool.to_string()), + Some(reasoning.to_string()), + ), + "vllm", + ) +} + #[tokio::test] async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { - let processor = |tool: &str, reasoning: &str| { - StreamingProcessor::new( - ToolParserFactory::new(), - ReasoningParserFactory::new(), - utils::ParserResolver::new( - Arc::new(WorkerRegistry::new()), - Some(tool.to_string()), - Some(reasoning.to_string()), - ), - "vllm", - ) - }; // `` ends reasoning without text, and the arguments of a specific // tool follow in the next chunk. let tool = messages::ToolChoice::Tool { @@ -798,7 +835,7 @@ async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { }; assert_eq!( messages_blocks( - processor("deepseek", "deepseek_r1"), + named_processor("deepseek", "deepseek_r1"), Some(tool), &["plan", "", "{}"], ) @@ -813,7 +850,7 @@ async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { // A reasoning parser that enters reasoning again after text. assert_eq!( messages_blocks( - processor("deepseek_v41", "deepseek_v41"), + named_processor("deepseek_v41", "deepseek_v41"), None, &[ "plan", @@ -838,3 +875,63 @@ async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { ] ); } + +#[tokio::test] +async fn messages_tool_arguments_precede_text_in_the_same_chunk() { + // With multi-token chunks, qwen_xml returns the arguments that finish a + // call together with the text after the call. + let (blocks, inputs) = messages_blocks_and_inputs( + named_processor("qwen_xml", "qwen3"), + None, + &[ + "\n\n\n1", + "\n\n\n\n\n\n", + "\n2\n\n\n", + ], + ) + .await; + assert_eq!( + blocks, + [ + "start 0 \"tool_use\"", + "stop 0", + "start 1 \"text\"", + "stop 1", + "start 2 \"tool_use\"", + "stop 2", + ] + ); + assert_eq!( + inputs, + [serde_json::json!({"q": 1}), serde_json::json!({"q": 2})] + ); + for (texts, input) in [ + ( + [ + "\n\n\nPar", + "is\n\n\n\nDone.", + ], + serde_json::json!({"q": "Paris"}), + ), + ( + [ + "\n\n\nx\n\n", + "\ny\n\n\n\n", + ], + serde_json::json!({"a": "x", "b": "y"}), + ), + ] { + let (blocks, inputs) = + messages_blocks_and_inputs(named_processor("qwen_xml", "qwen3"), None, &texts).await; + assert_eq!( + blocks, + [ + "start 0 \"tool_use\"", + "stop 0", + "start 1 \"text\"", + "stop 1" + ] + ); + assert_eq!(inputs, [input]); + } +} From 2f6be0b7ab8ccdd2bb7b59c321b5d191cb1063ed Mon Sep 17 00:00:00 2001 From: yechank <161688079+yechank-nvidia@users.noreply.github.com> Date: Wed, 30 Sep 2026 08:29:46 +0000 Subject: [PATCH 3/5] fix(grpc): drop specific tool arguments without an open Messages block With a specific tool_choice, the tool_use block starts at the first chunk outside reasoning. If reasoning starts only after that, as when `` is split across chunks, the thinking block stops the tool_use block, and the arguments that follow went to an index that was never started. The official Python SDK raises IndexError on such a delta. Send the arguments only while the tool_use block is open, and log and drop them otherwise. The client gets the tool_use block with empty input, as it did before the block stops were added. Signed-off-by: yechank <161688079+yechank-nvidia@users.noreply.github.com> --- .../src/routers/grpc/regular/streaming.rs | 26 +++++++++++-------- .../grpc/regular/streaming/eof_tests.rs | 26 +++++++++++++++++++ 2 files changed, 41 insertions(+), 11 deletions(-) diff --git a/model_gateway/src/routers/grpc/regular/streaming.rs b/model_gateway/src/routers/grpc/regular/streaming.rs index 7c3e98e5e..b23a46d9a 100644 --- a/model_gateway/src/routers/grpc/regular/streaming.rs +++ b/model_gateway/src/routers/grpc/regular/streaming.rs @@ -2394,19 +2394,23 @@ impl StreamingProcessor { .await?; tool_block_open = true; } - // Emit arguments delta + // Emit arguments delta, unless reasoning stopped the block if !normal_text.is_empty() { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockDelta { - index: current_block_index, - delta: ContentBlockDelta::InputJsonDelta { - partial_json: normal_text, + if tool_block_open { + Self::send_messages_event( + tx, + &mut sse_buffer, + &MessageStreamEvent::ContentBlockDelta { + index: current_block_index, + delta: ContentBlockDelta::InputJsonDelta { + partial_json: normal_text, + }, }, - }, - ) - .await?; + ) + .await?; + } else { + warn!("Dropping tool arguments without an open tool_use block"); + } } } else if let Some(ref mut parser) = streaming_tool_parser { // Regular/required tool choice: use incremental parser diff --git a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs index 945a6cc80..6c8743b5d 100644 --- a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs +++ b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs @@ -935,3 +935,29 @@ async fn messages_tool_arguments_precede_text_in_the_same_chunk() { assert_eq!(inputs, [input]); } } + +#[tokio::test] +async fn messages_specific_tool_arguments_need_an_open_block() { + // Reasoning that starts only after the specific tool's block stops that + // block; the arguments after it have no block to go to and are dropped. + let tool = messages::ToolChoice::Tool { + name: "lookup".to_string(), + disable_parallel_tool_use: None, + }; + let (blocks, inputs) = messages_blocks_and_inputs( + named_processor("qwen", "qwen3"), + Some(tool), + &["plan", "", "{}"], + ) + .await; + assert_eq!( + blocks, + [ + "start 0 \"tool_use\"", + "stop 0", + "start 1 \"thinking\"", + "stop 1", + ] + ); + assert_eq!(inputs, [serde_json::json!({})]); +} From 0a70316bc7741063456ea69b1f82f754d3669897 Mon Sep 17 00:00:00 2001 From: yechank <161688079+yechank-nvidia@users.noreply.github.com> Date: Wed, 30 Sep 2026 08:29:46 +0000 Subject: [PATCH 4/5] test(grpc): cover the remaining Messages block transitions Cover a thinking block that follows a tool_use block, and text that the tool parser releases at the end of the stream after reasoning started again. Replace the stub specific-tool case, which the deepseek_r1 case already covers, and correct two test comments. Signed-off-by: yechank <161688079+yechank-nvidia@users.noreply.github.com> --- .../grpc/regular/streaming/eof_tests.rs | 37 ++++++++++++++----- 1 file changed, 28 insertions(+), 9 deletions(-) diff --git a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs index 6c8743b5d..21db18d6c 100644 --- a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs +++ b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs @@ -616,7 +616,8 @@ impl ReasoningParser for ReasoningButText { fn mark_think_start_stripped(&mut self) {} } -/// Reports one whole call on its first chunk, before the chunk's text. +/// Returns the chunk as normal text and reports one whole `lookup` call on +/// its first parse. #[derive(Default)] struct CallFirst { called: bool, @@ -782,8 +783,10 @@ async fn messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternate( "stop 3", ] ); - // Text the reasoning parser returns while it stays in reasoning, as when - // one chunk ends reasoning, adds text and starts reasoning again. + // Text the reasoning parser returns while it stays in reasoning. Such text + // skips the tool parser and follows the chunk's reasoning, so a chunk like + // `answermore` still comes out in the wrong order; only the + // block boundaries are checked here. assert_eq!( messages_blocks(stub_processor(true), None, &["a", "text", "b"]).await, [ @@ -795,18 +798,16 @@ async fn messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternate( "stop 2", ] ); - // A specific tool right after reasoning. - let tool = messages::ToolChoice::Tool { - name: "lookup".to_string(), - disable_parallel_tool_use: None, - }; + // Reasoning right after a call. assert_eq!( - messages_blocks(stub_processor(false), Some(tool), &["a"]).await, + messages_blocks(stub_processor(false), None, &["a", "b"]).await, [ "start 0 \"thinking\"", "stop 0", "start 1 \"tool_use\"", "stop 1", + "start 2 \"thinking\"", + "stop 2", ] ); } @@ -874,6 +875,24 @@ async fn messages_blocks_do_not_overlap_with_deepseek_parsers() { "stop 3", ] ); + // Reasoning again while the tool parser holds text that it releases at + // the end of the stream. + assert_eq!( + messages_blocks( + named_processor("json", "deepseek_v41"), + None, + &["plan", "", "{", "", "more"], + ) + .await, + [ + "start 0 \"thinking\"", + "stop 0", + "start 1 \"thinking\"", + "stop 1", + "start 2 \"text\"", + "stop 2", + ] + ); } #[tokio::test] From 00e8603bce33136567f142c3d7b2f51355a0d8b3 Mon Sep 17 00:00:00 2001 From: yechank <161688079+yechank-nvidia@users.noreply.github.com> Date: Wed, 30 Sep 2026 11:05:22 +0000 Subject: [PATCH 5/5] fix(grpc): drop parsed tool arguments without an open Messages block Until it strips a ``, a reasoning parser such as qwen3 takes the first `` anywhere in the output as the start of reasoning. When the output does not open with ``, one inside a tool argument starts a thinking block in the middle of the call, and that block stops the tool_use block. The rest of the call's arguments, which the tool parser returns without a name, went to an index that was never started, or at the end of the stream to the open thinking block. The official Python SDK raises IndexError on the former. Send tool arguments through one helper that drops them when no tool_use block is open, as the specific tool path already did. The client gets the tool_use block without those arguments. On main the thinking block started on top of the tool_use block, so the part of the last argument that JSON-style parsers release at the end of the stream still reached the client there; it is dropped now. Log the drop at debug level, since it repeats for every chunk of the call. The open-block test now covers both cases. Signed-off-by: yechank <161688079+yechank-nvidia@users.noreply.github.com> --- .../src/routers/grpc/regular/streaming.rs | 112 +++++++++--------- .../grpc/regular/streaming/eof_tests.rs | 44 ++++++- 2 files changed, 96 insertions(+), 60 deletions(-) diff --git a/model_gateway/src/routers/grpc/regular/streaming.rs b/model_gateway/src/routers/grpc/regular/streaming.rs index b23a46d9a..85151d141 100644 --- a/model_gateway/src/routers/grpc/regular/streaming.rs +++ b/model_gateway/src/routers/grpc/regular/streaming.rs @@ -1842,6 +1842,30 @@ impl StreamingProcessor { Ok(()) } + /// Send tool call arguments to the open `tool_use` block. Reasoning can + /// stop that block in the middle of a call, and the arguments after it + /// have no block to go to, so they are dropped. + async fn send_tool_arguments( + tx: &SseSender, + buffer: &mut Vec, + index: u32, + tool_block_open: bool, + partial_json: String, + ) -> Result<(), String> { + if partial_json.is_empty() { + return Ok(()); + } + if !tool_block_open { + debug!("Dropping tool arguments without an open tool_use block"); + return Ok(()); + } + let delta = MessageStreamEvent::ContentBlockDelta { + index, + delta: ContentBlockDelta::InputJsonDelta { partial_json }, + }; + Self::send_messages_event(tx, buffer, &delta).await + } + /// Process reasoning content in Messages streaming mode (n=1 only). /// /// Returns `(normal_text, reasoning_text, in_reasoning)`. @@ -2395,23 +2419,14 @@ impl StreamingProcessor { tool_block_open = true; } // Emit arguments delta, unless reasoning stopped the block - if !normal_text.is_empty() { - if tool_block_open { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockDelta { - index: current_block_index, - delta: ContentBlockDelta::InputJsonDelta { - partial_json: normal_text, - }, - }, - ) - .await?; - } else { - warn!("Dropping tool arguments without an open tool_use block"); - } - } + Self::send_tool_arguments( + tx, + &mut sse_buffer, + current_block_index, + tool_block_open, + normal_text, + ) + .await?; } else if let Some(ref mut parser) = streaming_tool_parser { // Regular/required tool choice: use incremental parser match parser.parse_incremental(&normal_text, chat_tools).await { @@ -2430,19 +2445,14 @@ impl StreamingProcessor { 0 }; for tool_call_item in calls.drain(..finishing) { - if !tool_call_item.parameters.is_empty() { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockDelta { - index: current_block_index, - delta: ContentBlockDelta::InputJsonDelta { - partial_json: tool_call_item.parameters, - }, - }, - ) - .await?; - } + Self::send_tool_arguments( + tx, + &mut sse_buffer, + current_block_index, + tool_block_open, + tool_call_item.parameters, + ) + .await?; } // Emit normal text from parser as text content blocks @@ -2527,19 +2537,14 @@ impl StreamingProcessor { } // Emit incremental arguments - if !tool_call_item.parameters.is_empty() { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockDelta { - index: current_block_index, - delta: ContentBlockDelta::InputJsonDelta { - partial_json: tool_call_item.parameters, - }, - }, - ) - .await?; - } + Self::send_tool_arguments( + tx, + &mut sse_buffer, + current_block_index, + tool_block_open, + tool_call_item.parameters, + ) + .await?; } } Err(e) => { @@ -2678,19 +2683,14 @@ impl StreamingProcessor { tool_block_open = true; } - if !tool_call_item.parameters.is_empty() { - Self::send_messages_event( - tx, - &mut sse_buffer, - &MessageStreamEvent::ContentBlockDelta { - index: current_block_index, - delta: ContentBlockDelta::InputJsonDelta { - partial_json: tool_call_item.parameters, - }, - }, - ) - .await?; - } + Self::send_tool_arguments( + tx, + &mut sse_buffer, + current_block_index, + tool_block_open, + tool_call_item.parameters, + ) + .await?; } } } diff --git a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs index 21db18d6c..5f31623d3 100644 --- a/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs +++ b/model_gateway/src/routers/grpc/regular/streaming/eof_tests.rs @@ -686,8 +686,9 @@ async fn messages_blocks( /// The content block starts and stops of a Messages stream of `texts`, and /// the input of each `tool_use` block joined from its `input_json_delta`s as -/// a client SDK builds it. A block starts only when none is open, and every -/// delta goes to the open block and matches its type. +/// a client SDK builds it (the joined text if it is not complete JSON). A +/// block starts only when none is open, and every delta goes to the open +/// block and matches its type. async fn messages_blocks_and_inputs( processor: StreamingProcessor, tool_choice: Option, @@ -761,7 +762,7 @@ async fn messages_blocks_and_inputs( .iter() .map(|json| match json.as_str() { "" => serde_json::json!({}), - json => serde_json::from_str(json).expect("tool input JSON"), + json => serde_json::from_str(json).unwrap_or_else(|_| json.into()), }) .collect(); (blocks, inputs) @@ -956,7 +957,7 @@ async fn messages_tool_arguments_precede_text_in_the_same_chunk() { } #[tokio::test] -async fn messages_specific_tool_arguments_need_an_open_block() { +async fn messages_tool_arguments_need_an_open_block() { // Reasoning that starts only after the specific tool's block stops that // block; the arguments after it have no block to go to and are dropped. let tool = messages::ToolChoice::Tool { @@ -979,4 +980,39 @@ async fn messages_specific_tool_arguments_need_an_open_block() { ] ); assert_eq!(inputs, [serde_json::json!({})]); + // Until qwen3 strips a ``, it takes the first one anywhere as the + // start of reasoning, so one inside an argument stops the call's block. + // The rest of the arguments are dropped, in the middle of the stream and + // at its end, where qwen_xml releases the closing brace. + for (texts, input) in [ + ( + [ + "\n\n\nA ", + "x B\n\n", + "\n", + ], + serde_json::json!({}), + ), + ( + [ + "\n\n\nA\n\n", + "\nB ", + "x", + ], + serde_json::json!("{\"q\": \"A\""), + ), + ] { + let (blocks, inputs) = + messages_blocks_and_inputs(named_processor("qwen_xml", "qwen3"), None, &texts).await; + assert_eq!( + blocks, + [ + "start 0 \"tool_use\"", + "stop 0", + "start 1 \"thinking\"", + "stop 1", + ] + ); + assert_eq!(inputs, [input]); + } }