From 1a3ff347dc3ed12dbc94b64d1dec21a50f99146f Mon Sep 17 00:00:00 2001 From: ai-jz Date: Wed, 9 Sep 2026 15:57:05 -0700 Subject: [PATCH] fix(tool_parser): preserve Qwen XML streamed arguments Bound parameter extraction to the current tool call, close streamed root objects at tool boundaries, and recover a safe final brace at EOS. Add focused regressions for literal braces, truncation, and adjacent calls. Signed-off-by: ai-jz --- crates/tool_parser/src/parsers/qwen_xml.rs | 61 +++++--- .../tool_parser/tests/tool_parser_qwen_xml.rs | 145 ++++++++++++++++++ 2 files changed, 186 insertions(+), 20 deletions(-) diff --git a/crates/tool_parser/src/parsers/qwen_xml.rs b/crates/tool_parser/src/parsers/qwen_xml.rs index 37f0ea2b88..06f2c93dc8 100644 --- a/crates/tool_parser/src/parsers/qwen_xml.rs +++ b/crates/tool_parser/src/parsers/qwen_xml.rs @@ -174,13 +174,20 @@ impl QwenXmlParser { /// Parse and stream complete parameters from buffer /// Returns tool call items to emit (similar to Python's _parse_and_stream_parameters) - fn parse_and_stream_parameters(&mut self, tools: &[Tool]) -> Vec { + fn parse_and_stream_parameters( + &mut self, + tools: &[Tool], + parameter_end: usize, + ) -> Vec { let mut calls: Vec = vec![]; let param_types = helpers::param_types_for_function(tools, &self.current_function_name); - // Find all complete parameter patterns in buffer + // Leave parameters from subsequent coalesced calls for their own iteration. let mut new_params = serde_json::Map::new(); - for cap in self.xml_param_pattern.captures_iter(&self.buffer) { + for cap in self + .xml_param_pattern + .captures_iter(&self.buffer[..parameter_end]) + { if let (Some(key_match), Some(value_match)) = (cap.get(1), cap.get(2)) { let key = key_match.as_str().trim().to_string(); let value = value_match.as_str(); @@ -422,26 +429,23 @@ impl ToolParser for QwenXmlParser { // Parse parameters (only complete ones) if self.current_tool_name_sent { - let param_calls = self.parse_and_stream_parameters(tools); + let end_pos = self.buffer.find(self.tool_call_end_token); + let parameter_end = end_pos.unwrap_or(self.buffer.len()); + let param_calls = self.parse_and_stream_parameters(tools, parameter_end); calls.extend(param_calls); // Check if tool call is complete - if let Some(end_pos) = self.buffer.find(self.tool_call_end_token) { - // Close JSON object if we have parameters - let current_args = &self.streamed_args_for_tool[self.current_tool_id as usize]; - if !current_args.is_empty() { - // Count braces to check if JSON is complete - let open_braces = current_args.matches('{').count(); - let close_braces = current_args.matches('}').count(); - if open_braces > close_braces { - calls.push(ToolCallItem { - tool_index: self.current_tool_id as usize, - name: None, - parameters: "}".to_string(), - }); - self.streamed_args_for_tool[self.current_tool_id as usize].push('}'); - } - } + if let Some(end_pos) = end_pos { + // Parameter fragments leave the root open; braces in values are data. + let current_args = + &mut self.streamed_args_for_tool[self.current_tool_id as usize]; + let closing = if current_args.is_empty() { "{}" } else { "}" }; + calls.push(ToolCallItem { + tool_index: self.current_tool_id as usize, + name: None, + parameters: closing.to_string(), + }); + current_args.push_str(closing); // Complete the tool call self.buffer = @@ -466,6 +470,23 @@ impl ToolParser for QwenXmlParser { } fn get_unstreamed_tool_args(&self) -> Option> { + let tool_index = self.prev_tool_call_arr.len().checked_sub(1)?; + let actual = self.streamed_args_for_tool.get(tool_index)?; + let expected = self.prev_tool_call_arr[tool_index].get("arguments")?; + // XML parameters use spaced JSON, unlike the generic compact-prefix recovery. + // Append only the root brace, and only for exactly the already parsed values. + if !actual.is_empty() + && serde_json::from_str::(&format!("{actual}}}")) + .ok() + .as_ref() + == Some(expected) + { + return Some(vec![ToolCallItem { + tool_index, + name: None, + parameters: "}".to_string(), + }]); + } helpers::get_unstreamed_args(&self.prev_tool_call_arr, &self.streamed_args_for_tool) } diff --git a/crates/tool_parser/tests/tool_parser_qwen_xml.rs b/crates/tool_parser/tests/tool_parser_qwen_xml.rs index 564316a31f..74fbce5baa 100644 --- a/crates/tool_parser/tests/tool_parser_qwen_xml.rs +++ b/crates/tool_parser/tests/tool_parser_qwen_xml.rs @@ -8,6 +8,151 @@ use common::{create_test_tools, streaming_helpers}; use serde_json::json; use tool_parser::{parsers::QwenXmlParser, traits::ToolParser}; +// Exercise the same terminal getters as the gRPC streaming caller. Fixtures are synthetic. +#[expect(clippy::unwrap_used, reason = "assertions in a shared test helper")] +async fn assert_streamed_arguments( + input: &str, + expected: &[(&str, serde_json::Value)], + complete: bool, +) { + let tools = create_test_tools(); + let mut feeds = vec![ + vec![input.to_string()], + input.chars().map(|c| c.to_string()).collect(), + ]; + // Also cover every possible two-chunk split, including coalesced calls. + for (split, _) in input.char_indices().skip(1) { + feeds.push(vec![input[..split].to_string(), input[split..].to_string()]); + } + for chunks in feeds { + let mut parser = QwenXmlParser::new(); + let mut deltas = Vec::new(); + for chunk in &chunks { + let result = parser.parse_incremental(chunk, &tools).await.unwrap(); + assert!(result.normal_text.is_empty()); + deltas.extend(result.calls); + } + assert!(parser.take_unstreamed_normal_text().is_empty()); + let pending = parser.get_unstreamed_tool_args(); + if complete { + assert!(pending.is_none(), "completed calls must already be closed"); + assert!(parser + .parse_incremental("", &tools) + .await + .unwrap() + .calls + .is_empty()); + } + deltas.extend(pending.unwrap_or_default()); + let names: Vec<_> = deltas + .iter() + .filter_map(|d| d.name.as_deref().map(|name| (d.tool_index, name))) + .collect(); + let expected_names: Vec<_> = expected + .iter() + .enumerate() + .map(|(index, (name, _))| (index, *name)) + .collect(); + assert_eq!(names, expected_names, "chunks: {chunks:?}"); + assert!(deltas.iter().all(|d| d.tool_index < expected.len())); + for (index, (_, expected_args)) in expected.iter().enumerate() { + let args: String = deltas + .iter() + .filter(|d| d.tool_index == index) + .map(|d| d.parameters.as_str()) + .collect(); + let parsed = serde_json::from_str::(&args); + assert_eq!( + parsed.as_ref().ok(), + Some(expected_args), + "args: {args}; chunks: {chunks:?}" + ); + } + parser.reset(); + assert!(parser.get_unstreamed_tool_args().is_none()); + } +} + +#[tokio::test] +async fn test_qwen_xml_streaming_braces_inside_string_do_not_close_object() { + let input = "echo '}'"; + assert_streamed_arguments(input, &[("get_weather", json!({"city": "echo '}'"}))], true).await; +} + +#[tokio::test] +async fn test_qwen_xml_eos_closes_only_complete_parameter_values() { + let input = "Tokyo"; + for suffix in ["", "cels"] { + assert_streamed_arguments( + &format!("{input}{suffix}"), + &[("get_weather", json!({"city": "Tokyo"}))], + false, + ) + .await; + } +} + +#[tokio::test] +async fn test_qwen_xml_eos_does_not_invent_unfinished_parameter() { + assert_streamed_arguments( + "Tok", + &[("get_weather", json!({}))], + false, + ) + .await; +} + +#[tokio::test] +async fn test_qwen_xml_eos_closes_last_of_multiple_calls() { + let input = concat!( + "Paris", + "Tokyo", + ); + assert_streamed_arguments( + input, + &[ + ("get_weather", json!({"city": "Paris"})), + ("get_weather", json!({"city": "Tokyo"})), + ], + false, + ) + .await; +} + +#[tokio::test] +async fn test_qwen_xml_coalesced_calls_do_not_share_parameters() { + let input = concat!( + "Paris", + "Tokyocelsius", + ); + assert_streamed_arguments( + input, + &[ + ("get_weather", json!({"city": "Paris"})), + ("get_weather", json!({"city": "Tokyo", "units": "celsius"})), + ], + true, + ) + .await; +} + +#[tokio::test] +async fn test_qwen_xml_empty_calls_and_nested_values_are_closed_once() { + let input = concat!( + "", + r#"{"x":"}"}"#, + ); + assert_streamed_arguments( + input, + &[ + ("get_time", json!({})), + ("process", json!({"data": {"x": "}"}})), + ], + true, + ) + .await; +} + #[tokio::test] async fn test_qwen_xml_single_tool() { let parser = QwenXmlParser::new();