Conversation
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 `</think>` (deepseek_r1 with the deepseek tool parser, or qwen3 with qwen when the output opens with `<think>`), 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>
|
Navigate logical layers of code changes, visualize relationships, and explore their blast radius. No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Advanced Run ID: 📒 Files selected for processing (2)
Included review availability: This review used your included allowance. Your plan provides up to 4 included reviews per hour; 3 remain after this review. 📝 SummarySummary by CodeRabbit
WalkthroughMessages streaming now closes an open reasoning, text, or tool block before starting another. Tool arguments are emitted only when a tool block is open. Tests cover block transitions, argument ordering, and parser-flush output. ChangesMessages content block boundaries
Priority: ➖ Normal Estimated code review effort: 3 (Moderate) | ~20 minutes Change: Bug fix Merge Risk: ⚪ Minimal · up to The change closes open Messages blocks before transitions and keeps tool arguments within active tool blocks. No merge-blocking issue is established; merge after normal checks pass. Security Architecture ReviewSecurity architecture risk: 🔵 Low · up to The change improves stream ordering without adding tool-execution authority or shared request state. Interrupted tool calls can still finish with empty or incomplete arguments; whether downstream applications reject those calls remains unverified. Retained concerns Security review detailsSecurity Blast Radius
Trust Boundaries and Controls
Resilience and Maintainability Implications
Hardening Proposals
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1🛠️ Fix failing CI checks 💡
🧪 Generate unit tests (beta)
Comment |
…reaming 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</parameter>\n</function>\n</tool_call>\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>
With a specific tool_choice, the tool_use block starts at the first chunk outside reasoning. If reasoning starts only after that, as when `<think>` 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>
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>
Until it strips a `<think>`, a reasoning parser such as qwen3 takes the first `<think>` anywhere in the output as the start of reasoning. When the output does not open with `<think>`, 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>
Description
Problem
In short. In Messages streaming a response is a sequence of content blocks. Each block is sent as
content_block_start, its deltas, thencontent_block_stop, and the next block starts only after the previous one stops; clients buildmessage.contentfrom the block indexes. The gRPC path could start the next block without stopping the open one, sometimes at the same index (the specific-tool case below):Most streams still came out right: a response with a single block, or reasoning followed by text in a later chunk, never overlaps, and a client that only appends deltas by index rebuilds the content anyway. Where blocks overlap or an index is reused, deltas can reach the wrong block or a block that was never started: depending on the client the stream fails or a tool call arrives without its arguments, and a client that tracks a single open block sees blocks that never stop.
gRPC Messages streaming (
process_messages_streaming_chunks) keeps one open flag per block kind (thinking_block_open,text_block_open,tool_block_open) and onecurrent_block_index. Before this change:tool_usestart sites stopped an open block first, and only a text ortool_useblock;A block could then start at the index of a block that was still open. The deltas and stops that followed went to an index that was never started. Three cases reach this with the parsers on
main.Text after a tool call.
qwen_xmlreturns what follows</tool_call>, even a lone newline, as normal text. With one character per chunk,<tool_call>\n<function=lookup>\n<parameter=q>\n1\n</parameter>\n</function>\n</tool_call>\nDone.gives:A specific tool right after reasoning. Parsers
deepseek_r1+deepseek(orqwen3+qwenwhen the output opens with<think>),tool_choice: {"type": "tool", "name": "lookup"}, chunks["plan", "</think>", "{}"]:The
</think>chunk carries no text, so the thinking block is still open when the specific-tool path starts thetool_useblock. This needs a backend that lets the model reason before the tool's JSON schema applies.Reasoning again after text or a tool call.
deepseek_v41enters reasoning again on<think>in content.inklingcan send a thinking message after a text or tool message. Withdeepseek_v41and chunks["<think>plan", "</think>", "answer", "<think>", "more", "</think>", "done"]:Solution
Before every
content_block_start, stop whichever block is open and move to the next index. A new helper,stop_open_block, does this and replaces the existing text/tool_usestops at the threetool_usestart sites.Two changes keep tool arguments in their block now that a new block stops the open one:
qwen_xmldoes this when one chunk holds the end of a parameter,</tool_call>and the text after it. The router emitted the text first. With the stops above alone, the text block stopped thetool_useblock, the arguments went to the text block, and clients built the tool input without them (onmain, the text block started on top of thetool_useblock instead). The calls before the first one with a name continue the open call, and now go to the opentool_useblock before the text. The text in such a result follows the end of the call, so this keeps the model's order.tool_useblock is open, and are dropped otherwise. A thinking block can stop thetool_useblock in the middle of a call, and the arguments after it have no block to go to. See Out of scope for when this happens.A stream that never started a block while another was open emits the same events as before.
Changes
model_gateway/src/routers/grpc/regular/streaming.rsStreamingProcessor::stop_open_block.content_block_startsites inprocess_messages_streaming_chunks: thinking, text, text from the tool parser, the leftover text at end of stream, and thetool_usestarts of the specific-tool, incremental and end-of-stream paths.StreamingProcessor::send_tool_arguments, which sends a call's arguments only while atool_useblock is open, and drops them with a debug log otherwise. The specific-tool, incremental and end-of-stream paths send arguments through it.model_gateway/src/routers/grpc/regular/streaming/eof_tests.rsmessages_blocks_and_inputs: runs a scripted gRPC stream throughprocess_messages_streaming_chunks. It returns the block starts and stops, and the input of eachtool_useblock joined from itsinput_json_deltas as a client SDK builds it (the joined text if it is not complete JSON). It asserts that a block starts only when none is open, and that every delta goes to the open block and matches its type (text_deltato text,input_json_deltatotool_use,thinking_deltato thinking).messages_blocksreturns only the blocks.messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternate: stub parsers for three orders:messages_blocks_do_not_overlap_with_deepseek_parsers: the registered parsers from the cases above (deepseek_r1+deepseekwith a specific tool, anddeepseek_v41), andjson+deepseek_v41, where the tool parser releases held text at the end of the stream after reasoning started again.messages_tool_arguments_precede_text_in_the_same_chunk:qwen_xmlwith multi-token chunks: two parallel calls, arguments followed byDone., and a second parameter followed by the close and a newline. It checks the tool inputs.messages_tool_arguments_need_an_open_block:qwen+qwen3with a specific tool and<think>split across chunks, andqwen_xml+qwen3with<think>inside an argument, both in the middle of the stream and at its end (see Out of scope).The message of the first commit says that a thinking block was stopped only "when reasoning ended in a chunk that also carried text". The accurate statement is the second bullet under Problem: before the end of the stream, a thinking block was stopped only when a later chunk carried normal text outside reasoning. The commit is already pushed, so it is left as is.
Out of scope
A specific tool before reasoning. With a specific
tool_choice, thetool_useblock starts in the first chunk that is not in reasoning, even when that chunk has no text. If reasoning starts only after that, as when the first chunk holds a partial<think>, the thinking start now stops thetool_useblock, and the arguments after reasoning have no open block. They are dropped. For example,qwen3+qwenwith chunks["<thi", "nk>plan", "</think>", "{}"]gives:The client gets a
tool_useblock with empty input, as onmain. There, the thinking block starts on top of thetool_useblock (start 0 tool_use, start 0 thinking, ...), the arguments go to index 1, and the Python and TypeScript SDKs drop them. With the block stops alone, the arguments went to index 2, which was never started, and the Python SDK's stream accumulator raisedIndexErroron that delta. Fixing it means deciding when the specific-tool block should start, so it is left for a follow-up.<think>inside tool arguments. Until it strips a<think>, a reasoning parser such asqwen3takes the first<think>anywhere in the output as the start of reasoning. When the output does not open with<think>, as when thinking is enabled for a model that answers without it, a<think>inside a tool argument starts reasoning in the middle of the call. The thinking block now stops thetool_useblock, and the rest of the call's arguments are dropped.qwen3+qwen_xmlwith chunks["<tool_call>\n<function=lookup>\n<parameter=q>\nA ", "<think>x</think> B\n</parameter>\n", "</function>\n</tool_call>"]gives:The client gets a
tool_useblock with empty input, as onmain. There, the thinking block starts on top of thetool_useblock (start 0 tool_use, start 0 thinking, ...) and the arguments go to index 1. With the block stops alone, they went to index 2, which was never started, and the Python SDK raisedIndexError. When reasoning lasts to the end of the stream, what is dropped is the closing brace thatqwen_xmlreleases there, which with the block stops alone went to the open thinking block. The Python SDK parses the rest as partial JSON and builds the same input as onmain. Parsers that hold back the end of the last argument until the end of the stream, such asqwen,json,mistralandllama, lose that part too. Onmainit reached the client, because the thinking block had started at the index of thetool_useblock. Keeping these arguments would mean changing how the reasoning parser reads<think>, which this PR does not touch.Reasoning and text switching inside one chunk. The reasoning parser returns a chunk's reasoning and normal text as two strings, and the router emits the reasoning first. So when one chunk holds
</think>answer<think>more, the blocks no longer overlap, but they come out in the wrong order.deepseek_v41with chunks["<think>plan", "</think>answer<think>more", "rest</think>done"]gives:morejoins the first thinking block, and one reasoning span is split over two blocks. Normal text in a chunk that ends in reasoning also skips the tool parser and follows the chunk's reasoning. This is existing behavior, and the Chat Completions path does the same.Merging with #2538. #2538 adds
ReasoningParser::prompt_reasoningwithout a default, and the fieldMessagesResponseSpec::starts_in_reasoning. Whichever of the two lands second needsfn prompt_reasoning(&self, _: &str) -> PromptReasoning { PromptReasoning::Absent }in theReasoningButTexttest stub andstarts_in_reasoning: falsein theMessagesResponseSpecliteral ofmessages_blocks_and_inputs.Test Plan
Negative control. In scratch copies, I ran the new tests with the
streaming.rsofmain(61b250f7) and of the first commit (46448999). I also ran each case as its own test.main46448999messages_blocks_do_not_overlap_when_reasoning_calls_and_text_alternatemessages_blocks_do_not_overlap_with_deepseek_parsersmessages_tool_arguments_precede_text_in_the_same_chunkmessages_tool_arguments_need_an_open_blockThe cases added in this revision:
main46448999qwen_xml<tool_call>\n<function=lookup>\n<parameter=q>\n1,\n</parameter>\n</function>\n</tool_call>\n<tool_call>\n<function=lookup>\n,<parameter=q>\n2\n</parameter>\n</function>\n</tool_call>start 0 tool_use,start 0 text{}instead of{"q": 1}qwen_xml<tool_call>\n<function=lookup>\n<parameter=q>\nPar,is\n</parameter>\n</function>\n</tool_call>\nDone.start 0 tool_use,start 0 text{}instead of{"q": "Paris"}qwen_xml<tool_call>\n<function=lookup>\n<parameter=a>\nx\n</parameter>\n,<parameter=b>\ny\n</parameter>\n</function>\n</tool_call>\nstart 0 tool_use,start 0 text{"a": "x"}instead of{"a": "x", "b": "y"}<think>qwen3+qwen<thi,nk>plan,</think>,{}start 0 tool_use,start 0 thinkinginput_json_deltaat index 2, which was never started<think>inside an argumentqwen3+qwen_xml<tool_call>\n<function=lookup>\n<parameter=q>\nA,<think>x</think> B\n</parameter>\n,</function>\n</tool_call>start 0 tool_use,start 0 thinkinginput_json_deltaat index 2, which was never started<think>inside an argument, then the end of the streamqwen3+qwen_xml<tool_call>\n<function=lookup>\n<parameter=q>\nA\n</parameter>\n,<parameter=r>\nB,<think>xstart 0 tool_use,start 0 thinkinginput_json_deltato the open thinking blocka,bstart 0 thinking,start 0 tool_usejson+deepseek_v41<think>plan,</think>,{,<think>,morestart 0 thinking,stop 0,start 1 thinking,start 1 textqwen_xmlstart 0 tool_use,start 0 textThe other cases of the first revision (the stub specific-tool case was replaced) still fail on
mainand pass with this PR.Mutation check. In a scratch copy, I disabled one change at a time and ran
eof_tests(test names shortened):stop_open_blockbefore the thinking start..._alternate,..._deepseek_parsers,..._need_an_open_block..._alternate,..._need_an_open_blocktool_usestart..._deepseek_parsers..._alternate,..._precede_text_in_the_same_chunktool_usestart..._alternate,..._precede_text_in_the_same_chunk..._alternate..._deepseek_parserstool_usestart..._precede_text_in_the_same_chunksend_tool_arguments..._need_an_open_block..._need_an_open_block(each)Only a parser whose
get_unstreamed_tool_argsreturns a named call reaches the end-of-streamtool_usestart while another block is open. Every registered parser returns unnamed items there, so no test covers that call.Commands. All ran offline on CPU.
cargo +nightly fmt --all -- --checkcargo test -p smg --lib -- eof_testscargo test -p smg --libcargo test -p smg --testmessages_streaming_test/messages_test/grpc_responses_stream_contract_test/api_tests/spec_test/grpc_context_length_test/grpc_pd_fanout_testcargo clippy --workspace --all-targets -- -D warnings;-p smg --all-targetswith default features and withgrpc-server,jemalloc-profiling,test-util; CI's two--no-default-featuresvariantspre-commit run --all-fileswith CI'sSKIPlistThe offline crate cache did not have the three dependency bumps on
main:lru0.18.5,cc1.5.1 and theopentelemetry-proto0.33 dev-dependency. So the build copy used theCargo.lockfrom before those bumps and the dev-dependency at 0.32. No source file differed.Checklist
cargo +nightly fmtpassescargo clippy --all-targets -- -D warningspasses for the workspace,-p smgand CI's two--no-default-featuresvariants.--all-featuresneeds OpenCV and was not run offline.