Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 41 additions & 20 deletions crates/tool_parser/src/parsers/qwen_xml.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<ToolCallItem> {
fn parse_and_stream_parameters(
&mut self,
tools: &[Tool],
parameter_end: usize,
) -> Vec<ToolCallItem> {
let mut calls: Vec<ToolCallItem> = 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])

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: Truncating the scan at the first </tool_call> silently changes behavior for parameter values that contain the literal </tool_call>, and nothing pins the new behavior.

For <tool_call><function=process><parameter=snippet>use </tool_call> tag</parameter></function></tool_call>:

  • Before: the regex ran over the whole buffer, captured snippet = "use </tool_call> tag", and the brace counter appended } — the client received the complete, valid {"snippet": "use </tool_call> tag"} (plus the leftover tag</parameter>… leaking as normal text).
  • After: parameter_end lands mid-value, the slice is <parameter=snippet>use with no </parameter>, so the parameter never matches and the call closes as {} — the argument is dropped entirely.

This is arguably the right call: parse_complete_inner's non-greedy (?s)<tool_call>\s*(.*?)\s*</tool_call> extractor already truncates at the same point and yields {}, so streaming and non-streaming now agree, and vLLM/SGLang split the same way. But the PR description lists this under "existing delimiter ambiguity … remain outside this change", when the change actually flips it from "correct args" to "empty args". Worth a short regression test in tool_parser_qwen_xml.rs asserting json!({}) for that fixture, so the alignment with parse_complete is intentional and stays that way.

{
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();
Expand Down Expand Up @@ -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 =
Expand All @@ -466,6 +470,23 @@ impl ToolParser for QwenXmlParser {
}

fn get_unstreamed_tool_args(&self) -> Option<Vec<ToolCallItem>> {
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::<Value>(&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)
}

Expand Down
145 changes: 145 additions & 0 deletions crates/tool_parser/tests/tool_parser_qwen_xml.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<serde_json::Value>(&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 = "<tool_call><function=get_weather><parameter=city>echo '}'</parameter></function></tool_call>";
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 = "<tool_call><function=get_weather><parameter=city>Tokyo</parameter>";
for suffix in ["", "<parameter=units>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(
"<tool_call><function=get_weather><parameter=city>Tok",
&[("get_weather", json!({}))],
false,
)
.await;
}

#[tokio::test]
async fn test_qwen_xml_eos_closes_last_of_multiple_calls() {
let input = concat!(
"<tool_call><function=get_weather><parameter=city>Paris</parameter></function></tool_call>",
"<tool_call><function=get_weather><parameter=city>Tokyo</parameter>",
);
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!(
"<tool_call><function=get_weather><parameter=city>Paris</parameter></function></tool_call>",
"<tool_call><function=get_weather><parameter=city>Tokyo</parameter><parameter=units>celsius</parameter></function></tool_call>",
);
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!(
"<tool_call><function=get_time></function></tool_call>",
r#"<tool_call><function=process><parameter=data>{"x":"}"}</parameter></function></tool_call>"#,
);
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();
Expand Down
Loading