use std::collections::BTreeMap;
use std::fs;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use dynamo_llm::model_card::ModelDeploymentCard;
use dynamo_llm::preprocessor::OpenAIPreprocessor;
use dynamo_llm::protocols::openai::chat_completions::{
NvCreateChatCompletionRequest, NvCreateChatCompletionStreamResponse,
};
use dynamo_protocols::types::{
ChatCompletionMessageContent, ChatCompletionNamedToolChoice, ChatCompletionTool,
ChatCompletionToolChoiceOption, ChatCompletionToolType, FinishReason, FunctionName,
};
use dynamo_runtime::protocols::annotated::Annotated;
use futures::{StreamExt, stream};
use serde_json::Value;
const REQUEST_JSON: &str = r#"{"messages":[{"role":"user","content":"What is the capital of Tuvalu?"}],"model":"Qwen/Qwen3-0.6B","max_completion_tokens":3000,"stream":true,"stream_options":{"include_usage":true,"continuous_usage_stats":false},"temperature":1.0,"top_p":1.0}"#;
const FORCE_REASONING_PARSERS: &[&str] = &[
"deepseek_r1",
"deepseek_v3",
"deepseek_v3_1",
"deepseek_v3_2",
"step3",
"kimi_k25",
"mistral",
"minimax_append_think",
"nemotron_nano",
"nemotron3",
"nemotron_v3",
];
const REASONING_BEFORE_GUIDED_JSON_PARSERS: &[(&str, &str)] = &[
("deepseek_r1", "</think>"),
("deepseek_v3", "</think>"),
("deepseek_v3_1", "</think>"),
("deepseek_v3_2", "</think>"),
("step3", "</think>"),
("kimi_k25", "</think>"),
("mistral", "[/THINK]"),
("nemotron_nano", "</think>"),
("nemotron3", "</think>"),
("nemotron_v3", "</think>"),
];
fn build_preprocessor(
reasoning_parser: Option<&str>,
tool_call_parser: Option<&str>,
) -> Arc<OpenAIPreprocessor> {
let model_path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/data/sample-models/mock-llama-3.1-8b-instruct");
let mut mdc = ModelDeploymentCard::load_from_disk(model_path, None).unwrap();
mdc.runtime_config.reasoning_parser = reasoning_parser.map(ToString::to_string);
mdc.runtime_config.tool_call_parser = tool_call_parser.map(ToString::to_string);
OpenAIPreprocessor::new(mdc).unwrap()
}
fn fixture_path(name: &str) -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("tests/data/replays")
.join(name)
}
fn parse_fixture(
jsonl_path: &Path,
) -> (
NvCreateChatCompletionRequest,
Vec<Value>,
Vec<NvCreateChatCompletionStreamResponse>,
) {
let content = fs::read_to_string(jsonl_path)
.unwrap_or_else(|e| panic!("failed to read fixture {}: {e}", jsonl_path.display()));
let mut expected_stream_json = Vec::new();
let mut input_chunks = Vec::new();
for line in content.lines().filter(|l| !l.is_empty()) {
let value: Value = serde_json::from_str(line).unwrap();
let chunk: NvCreateChatCompletionStreamResponse =
serde_json::from_value(value.clone()).unwrap();
let normalized = serde_json::to_value(&chunk).unwrap();
expected_stream_json.push(normalized);
input_chunks.push(chunk);
}
let request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
assert!(
!input_chunks.is_empty(),
"missing stream chunks in fixture {}",
jsonl_path.display()
);
(request, expected_stream_json, input_chunks)
}
fn get_text(content: &ChatCompletionMessageContent) -> &str {
match content {
ChatCompletionMessageContent::Text(text) => text.as_str(),
ChatCompletionMessageContent::Parts(_) => "",
}
}
#[derive(Default, Clone, Debug)]
struct MergedToolCall {
id: Option<String>,
r#type: Option<String>,
name: Option<String>,
arguments: String,
}
impl MergedToolCall {
fn merge_from(
&mut self,
tool_call: &dynamo_protocols::types::ChatCompletionMessageToolCallChunk,
) {
if self.id.is_none() {
self.id = tool_call.id.clone();
}
if self.r#type.is_none() {
self.r#type = tool_call.r#type.as_ref().map(|t| {
serde_json::to_string(t)
.unwrap()
.trim_matches('"')
.to_string()
});
}
if let Some(function) = &tool_call.function {
if self.name.is_none() {
self.name = function.name.clone();
}
if let Some(arguments) = &function.arguments {
self.arguments.push_str(arguments);
}
}
}
}
#[tokio::test]
async fn postprocessor_parsing_stream_replays_unit_test_fixture() {
let preprocessor = build_preprocessor(None, None);
let (request, expected_stream_json, input_chunks) =
parse_fixture(&fixture_path("stream_interval_1.jsonl"));
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
assert_eq!(output_chunks.len(), expected_stream_json.len());
for (idx, (output, expected)) in output_chunks
.iter()
.zip(expected_stream_json.iter())
.enumerate()
{
let output_data = output
.data
.as_ref()
.expect("output stream chunk should include data");
let output_json = serde_json::to_value(output_data).unwrap();
assert_eq!(output_json, *expected, "chunk {idx} did not match fixture");
}
}
#[tokio::test]
async fn postprocessor_parsing_stream_replays_interval_20_fixture() {
let preprocessor = build_preprocessor(Some("qwen"), Some("hermes"));
let (mut request, _expected_stream_json, input_chunks) =
parse_fixture(&fixture_path("stream_interval_20.jsonl"));
let tools: Vec<dynamo_protocols::types::ChatCompletionTool> =
serde_json::from_value(serde_json::json!([
{
"type": "function",
"function": {
"name": "search_gutenberg_books",
"description": "Search for books in the Project Gutenberg library",
"parameters": {
"type": "object",
"properties": {
"search_terms": {
"type": "array",
"items": {"type": "string"},
"description": "List of search terms to find books"
}
},
"required": ["search_terms"]
}
}
}
]))
.unwrap();
request.inner.tools = Some(tools);
request.inner.tool_choice = Some(ChatCompletionToolChoiceOption::Auto);
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut all_content = String::new();
let mut finish_reasons = Vec::new();
let mut merged_tool_calls: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
for output in &output_chunks {
let Some(output_data) = output.data.as_ref() else {
continue;
};
for choice in &output_data.inner.choices {
if let Some(reasoning_content) = &choice.delta.reasoning_content {
reasoning.push_str(reasoning_content);
}
if let Some(content) = &choice.delta.content {
all_content.push_str(get_text(content));
}
if let Some(reason) = choice.finish_reason {
finish_reasons.push(reason);
}
if let Some(tool_calls) = &choice.delta.tool_calls {
for tool_call in tool_calls {
merged_tool_calls
.entry(tool_call.index)
.or_default()
.merge_from(tool_call);
}
}
}
}
let tool_calls: Vec<MergedToolCall> = merged_tool_calls.values().cloned().collect();
assert!(
reasoning.contains("the user is asking for the titles of some James Joyce books"),
"reasoning did not contain expected phrase: {reasoning}"
);
assert!(
reasoning.contains("the user's request.\n"),
"reasoning did not contain expected ending: {reasoning}"
);
assert_eq!(
tool_calls.len(),
1,
"Expected 1 tool call but got {}. Tool-call markup was likely emitted as plain content instead.",
tool_calls.len()
);
let tc = &tool_calls[0];
assert_eq!(tc.name.as_deref(), Some("search_gutenberg_books"));
let arguments_json: Value = serde_json::from_str(&tc.arguments).unwrap();
assert_eq!(
arguments_json,
serde_json::json!({
"search_terms": ["James Joyce", "Project Gutenberg"]
})
);
assert!(
tc.id
.as_ref()
.is_some_and(|id| id.starts_with("call-") || id.starts_with("chatcmpl-tool-")),
"tool call id did not match expected prefix: {:?}",
tc.id
);
assert_eq!(tc.r#type.as_deref(), Some("function"));
assert!(
!all_content.contains("<tool_call>"),
"Raw <tool_call> markup leaked into content: {all_content:?}"
);
assert!(!all_content.contains("</tool_call>"));
if !finish_reasons.is_empty() {
assert!(
finish_reasons.contains(&FinishReason::Stop)
|| finish_reasons.contains(&FinishReason::ToolCalls),
"expected terminal finish reason (stop/tool_calls), got: {:?}",
finish_reasons
);
}
}
fn mock_content_chunk(content: &str) -> NvCreateChatCompletionStreamResponse {
use dynamo_protocols::types::{
ChatChoiceStream, ChatCompletionStreamResponseDelta, CreateChatCompletionStreamResponse,
Role,
};
#[allow(deprecated)]
let choice = ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
role: Some(Role::Assistant),
content: Some(ChatCompletionMessageContent::Text(content.to_string())),
tool_calls: None,
function_call: None,
refusal: None,
reasoning_content: None,
},
finish_reason: None,
logprobs: None,
};
NvCreateChatCompletionStreamResponse {
inner: CreateChatCompletionStreamResponse {
id: "test-id".to_string(),
choices: vec![choice],
created: 0,
model: "test-model".to_string(),
system_fingerprint: None,
object: "chat.completion.chunk".to_string(),
usage: None,
service_tier: None,
},
nvext: None,
llm_metrics: None,
}
}
fn mock_multi_choice_content_chunk(
choices: &[(u32, &str)],
) -> NvCreateChatCompletionStreamResponse {
use dynamo_protocols::types::{
ChatChoiceStream, ChatCompletionStreamResponseDelta, CreateChatCompletionStreamResponse,
Role,
};
#[allow(deprecated)]
let choices = choices
.iter()
.map(|(index, content)| ChatChoiceStream {
index: *index,
delta: ChatCompletionStreamResponseDelta {
role: Some(Role::Assistant),
content: Some(ChatCompletionMessageContent::Text((*content).to_string())),
tool_calls: None,
function_call: None,
refusal: None,
reasoning_content: None,
},
finish_reason: None,
logprobs: None,
})
.collect();
NvCreateChatCompletionStreamResponse {
inner: CreateChatCompletionStreamResponse {
id: "test-id".to_string(),
choices,
created: 0,
model: "test-model".to_string(),
system_fingerprint: None,
object: "chat.completion.chunk".to_string(),
usage: None,
service_tier: None,
},
nvext: None,
llm_metrics: None,
}
}
fn mock_reasoning_only_chunk(reasoning: &str) -> NvCreateChatCompletionStreamResponse {
use dynamo_protocols::types::{
ChatChoiceStream, ChatCompletionStreamResponseDelta, CreateChatCompletionStreamResponse,
Role,
};
#[allow(deprecated)]
let choice = ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
role: Some(Role::Assistant),
content: None,
tool_calls: None,
function_call: None,
refusal: None,
reasoning_content: Some(reasoning.to_string()),
},
finish_reason: None,
logprobs: None,
};
NvCreateChatCompletionStreamResponse {
inner: CreateChatCompletionStreamResponse {
id: "test-id".to_string(),
choices: vec![choice],
created: 0,
model: "test-model".to_string(),
system_fingerprint: None,
object: "chat.completion.chunk".to_string(),
usage: None,
service_tier: None,
},
nvext: None,
llm_metrics: None,
}
}
fn mock_final_chunk() -> NvCreateChatCompletionStreamResponse {
use dynamo_protocols::types::{
ChatChoiceStream, ChatCompletionStreamResponseDelta, CreateChatCompletionStreamResponse,
};
#[allow(deprecated)]
let choice = ChatChoiceStream {
index: 0,
delta: ChatCompletionStreamResponseDelta {
role: None,
content: None,
tool_calls: None,
function_call: None,
refusal: None,
reasoning_content: None,
},
finish_reason: Some(FinishReason::Stop),
logprobs: None,
};
NvCreateChatCompletionStreamResponse {
inner: CreateChatCompletionStreamResponse {
id: "test-id".to_string(),
choices: vec![choice],
created: 0,
model: "test-model".to_string(),
system_fingerprint: None,
object: "chat.completion.chunk".to_string(),
usage: None,
service_tier: None,
},
nvext: None,
llm_metrics: None,
}
}
#[tokio::test]
async fn postprocessor_parsing_stream_deepseek_v4_tool_continuation_keeps_injected_reasoning() {
let preprocessor = build_preprocessor(Some("deepseek_v4"), None);
let request: NvCreateChatCompletionRequest = serde_json::from_value(serde_json::json!({
"messages": [
{"role": "user", "content": "Create and run a hello-world script."},
{
"role": "assistant",
"tool_calls": [{
"id": "call_1",
"type": "function",
"function": {
"name": "run_python",
"arguments": "{\"path\":\"/tmp/hello.py\"}"
}
}]
},
{
"role": "tool",
"tool_call_id": "call_1",
"content": "Hello, world!"
}
],
"model": "deepseek-ai/DeepSeek-V4-Pro",
"stream": true
}))
.unwrap();
let input_chunks = vec![
mock_content_chunk("The script ran successfully."),
mock_content_chunk("</think>"),
mock_content_chunk("Done. Output: `Hello, world!`"),
mock_final_chunk(),
];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
}
}
assert_eq!(reasoning, "The script ran successfully.");
assert_eq!(content, "Done. Output: `Hello, world!`");
assert!(
!content.contains("</think>"),
"literal closing tag leaked into content: {content:?}"
);
}
fn kimi_tool_continuation_request(
model: &str,
thinking: Option<bool>,
) -> NvCreateChatCompletionRequest {
let mut request = serde_json::json!({
"messages": [
{"role": "user", "content": "What is the weather in London?"},
{
"role": "assistant",
"tool_calls": [{
"id": "call_1",
"type": "function",
"function": {
"name": "get_weather",
"arguments": "{\"location\":\"London\"}"
}
}]
},
{
"role": "tool",
"tool_call_id": "call_1",
"content": "{\"temperature\":15,\"unit\":\"celsius\",\"condition\":\"cloudy\"}"
}
],
"model": model,
"stream": true
});
if let Some(thinking) = thinking {
request["chat_template_kwargs"] = serde_json::json!({"thinking": thinking});
}
serde_json::from_value(request).unwrap()
}
async fn run_kimi_tool_continuation(
request: NvCreateChatCompletionRequest,
input_chunks: Vec<NvCreateChatCompletionStreamResponse>,
) -> DrainOutput {
let preprocessor = build_preprocessor(Some("kimi_k25"), None);
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
drain_stream(output_stream).await
}
#[tokio::test]
async fn postprocessor_parsing_stream_kimi_k25_tool_continuation_with_thinking_parses_reasoning() {
let request = kimi_tool_continuation_request("moonshotai/Kimi-K2.6", Some(true));
let output = run_kimi_tool_continuation(
request,
vec![
mock_content_chunk("The tool returned 15°C and cloudy."),
mock_content_chunk("</thi"),
mock_content_chunk("nk>The current weather in London is 15°C and cloudy."),
mock_final_chunk(),
],
)
.await;
assert_eq!(output.reasoning, "The tool returned 15°C and cloudy.");
assert_eq!(
output.content,
"The current weather in London is 15°C and cloudy."
);
assert!(
!output.content.contains("</think>"),
"literal closing tag leaked into content: {:?}",
output.content
);
}
#[tokio::test]
async fn postprocessor_parsing_stream_kimi_k26_omitted_thinking_parses_reasoning() {
let request = kimi_tool_continuation_request("moonshotai/Kimi-K2.6", None);
let output = run_kimi_tool_continuation(
request,
vec![
mock_content_chunk("The tool returned 15°C and cloudy."),
mock_content_chunk("</thi"),
mock_content_chunk("nk>The current weather in London is 15°C and cloudy."),
mock_final_chunk(),
],
)
.await;
assert_eq!(
(output.reasoning.as_str(), output.content.as_str()),
(
"The tool returned 15°C and cloudy.",
"The current weather in London is 15°C and cloudy."
)
);
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_v3_enable_thinking_false_returns_content() {
let preprocessor = build_preprocessor(Some("nemotron_v3"), None);
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
request.chat_template_args = Some(
serde_json::from_value(serde_json::json!({
"enable_thinking": false
}))
.unwrap(),
);
let input_chunks = vec![mock_content_chunk("This is plain content")];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
}
}
assert_eq!(reasoning, "");
assert_eq!(content, "This is plain content");
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_v3_force_nonempty_strips_start_token() {
let preprocessor = build_preprocessor(Some("nemotron_v3"), None);
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
request.chat_template_args = Some(
serde_json::from_value(serde_json::json!({
"force_nonempty_content": true
}))
.unwrap(),
);
let input_chunks = vec![
mock_content_chunk("<thi"),
mock_content_chunk("nk>This is plain content"),
];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
}
}
assert_eq!(reasoning, "");
assert_eq!(content, "This is plain content");
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_v3_force_nonempty_flushes_partial_prefix_on_finish()
{
let preprocessor = build_preprocessor(Some("nemotron_v3"), None);
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
request.chat_template_args = Some(
serde_json::from_value(serde_json::json!({
"force_nonempty_content": true
}))
.unwrap(),
);
let input_chunks = vec![mock_content_chunk("<thi"), mock_final_chunk()];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
let mut finish_reasons = Vec::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
if let Some(fr) = choice.finish_reason {
finish_reasons.push(fr);
}
}
}
assert_eq!(reasoning, "");
assert_eq!(content, "<thi");
assert!(finish_reasons.contains(&FinishReason::Stop));
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_v3_force_nonempty_flushes_partial_prefix_on_eof() {
let preprocessor = build_preprocessor(Some("nemotron_v3"), None);
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
request.chat_template_args = Some(
serde_json::from_value(serde_json::json!({
"force_nonempty_content": true
}))
.unwrap(),
);
let input_chunks = vec![mock_content_chunk("<thi")];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
}
}
assert_eq!(reasoning, "");
assert_eq!(content, "<thi");
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_v3_force_nonempty_tracks_prefix_per_choice() {
let preprocessor = build_preprocessor(Some("nemotron_v3"), None);
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
request.chat_template_args = Some(
serde_json::from_value(serde_json::json!({
"force_nonempty_content": true
}))
.unwrap(),
);
let input_chunks = vec![
mock_multi_choice_content_chunk(&[(0, "<thi"), (1, "<thi")]),
mock_multi_choice_content_chunk(&[(0, "nk>First"), (1, "nk>Second")]),
];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut content_by_choice = BTreeMap::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(c) = &choice.delta.content {
content_by_choice
.entry(choice.index)
.or_insert_with(String::new)
.push_str(get_text(c));
}
assert!(
choice.delta.reasoning_content.is_none(),
"reasoning_content must stay empty when force_nonempty_content=true"
);
}
}
assert_eq!(content_by_choice.get(&0).map(String::as_str), Some("First"));
assert_eq!(
content_by_choice.get(&1).map(String::as_str),
Some("Second")
);
}
#[tokio::test]
async fn postprocessor_parsing_stream_minimax_required_bypasses_reasoning() {
let preprocessor = build_preprocessor(Some("minimax_append_think"), Some("minimax_m2"));
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
let tools: Vec<dynamo_protocols::types::ChatCompletionTool> =
serde_json::from_value(serde_json::json!([{
"type": "function",
"function": {
"name": "get_weather",
"description": "Get the current weather for a location.",
"parameters": {
"type": "object",
"properties": {"location": {"type": "string"}},
"required": ["location"]
}
}
}]))
.unwrap();
request.inner.tools = Some(tools);
request.inner.tool_choice = Some(ChatCompletionToolChoiceOption::Required);
let bare_json = r#"[{"name": "get_weather", "parameters": {"location": "San Francisco"}}]"#;
let input_chunks = vec![mock_content_chunk(bare_json), mock_final_chunk()];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
let mut merged_tool_calls: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
let mut finish_reasons = Vec::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
if let Some(tcs) = &choice.delta.tool_calls {
for tc in tcs {
merged_tool_calls
.entry(tc.index)
.or_default()
.merge_from(tc);
}
}
if let Some(fr) = choice.finish_reason {
finish_reasons.push(fr);
}
}
}
assert!(
reasoning.is_empty(),
"reasoning_content must be empty when tool_choice=required forces bare JSON, got: {reasoning:?}"
);
assert!(
!content.contains("get_weather"),
"tool call JSON must not leak into content, got: {content:?}"
);
let tool_calls: Vec<MergedToolCall> = merged_tool_calls.values().cloned().collect();
assert_eq!(tool_calls.len(), 1, "expected one tool call");
assert_eq!(tool_calls[0].name.as_deref(), Some("get_weather"));
let args: Value = serde_json::from_str(&tool_calls[0].arguments).unwrap();
assert_eq!(args, serde_json::json!({"location": "San Francisco"}));
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn postprocessor_parsing_stream_nemotron_required_smoke_case() {
for (case, parser, prompt_injected_reasoning) in [
("nano", "nemotron_nano", true),
("super/deci", "nemotron_deci", false),
] {
let preprocessor = build_preprocessor(Some(parser), Some(parser));
let mut request: NvCreateChatCompletionRequest =
serde_json::from_value(serde_json::json!({
"model": "nvidia/nvidia/nemotron-3-super-120b-long-ctx",
"messages": [
{"role": "user", "content": "What is the weather in San Francisco?"}
],
"stream": true,
"temperature": 0.0
}))
.unwrap();
let tools: Vec<dynamo_protocols::types::ChatCompletionTool> =
serde_json::from_value(serde_json::json!([{
"type": "function",
"function": {
"name": "get_weather",
"description": "Get the current weather for a location.",
"parameters": {
"type": "object",
"properties": {
"location": {
"type": "string",
"description": "City name"
}
},
"required": ["location"]
}
}
}]))
.unwrap();
request.inner.tools = Some(tools);
request.inner.tool_choice = Some(ChatCompletionToolChoiceOption::Required);
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let input_chunks = vec![mock_content_chunk(bare_json), mock_final_chunk()];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, prompt_injected_reasoning, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
let mut merged_tool_calls: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
let mut finish_reasons = Vec::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
if let Some(tcs) = &choice.delta.tool_calls {
for tc in tcs {
merged_tool_calls
.entry(tc.index)
.or_default()
.merge_from(tc);
}
}
if let Some(fr) = choice.finish_reason {
finish_reasons.push(fr);
}
}
}
assert!(
reasoning.is_empty(),
"{case}: reasoning_content must be empty when tool_choice=required forces bare JSON, got: {reasoning:?}"
);
assert!(
!content.contains("get_weather"),
"{case}: tool-call JSON must not leak into content, got: {content:?}"
);
assert!(
!content.contains("<tool_call>"),
"{case}: raw <tool_call> XML must not leak into content, got: {content:?}"
);
let tool_calls: Vec<MergedToolCall> = merged_tool_calls.values().cloned().collect();
assert_eq!(tool_calls.len(), 1, "{case}: expected one tool call");
assert_eq!(
tool_calls[0].name.as_deref(),
Some("get_weather"),
"{case}: wrong tool name"
);
let args: Value = serde_json::from_str(&tool_calls[0].arguments).unwrap();
assert_eq!(
args,
serde_json::json!({"location": "San Francisco"}),
"{case}: wrong arguments"
);
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
#[tokio::test]
async fn postprocessor_parsing_stream_minimax_named_bypasses_reasoning() {
let preprocessor = build_preprocessor(Some("minimax_append_think"), Some("minimax_m2"));
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
let tools: Vec<dynamo_protocols::types::ChatCompletionTool> =
serde_json::from_value(serde_json::json!([{
"type": "function",
"function": {
"name": "get_weather",
"description": "Get the current weather for a location.",
"parameters": {
"type": "object",
"properties": {"location": {"type": "string"}},
"required": ["location"]
}
}
}]))
.unwrap();
request.inner.tools = Some(tools);
request.inner.tool_choice = Some(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
let bare_json = r#"[{"name": "get_weather", "parameters": {"location": "Tokyo"}}]"#;
let input_chunks = vec![mock_content_chunk(bare_json), mock_final_chunk()];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut merged_tool_calls: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
let mut finish_reasons = Vec::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(tcs) = &choice.delta.tool_calls {
for tc in tcs {
merged_tool_calls
.entry(tc.index)
.or_default()
.merge_from(tc);
}
}
if let Some(fr) = choice.finish_reason {
finish_reasons.push(fr);
}
}
}
assert!(
reasoning.is_empty(),
"reasoning_content must be empty for tool_choice=named, got: {reasoning:?}"
);
let tool_calls: Vec<MergedToolCall> = merged_tool_calls.values().cloned().collect();
assert_eq!(tool_calls.len(), 1);
assert_eq!(tool_calls[0].name.as_deref(), Some("get_weather"));
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"named tool_choice with emitted tool_calls should finish as ToolCalls, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn postprocessor_parsing_stream_minimax_named_bare_parameters() {
let preprocessor = build_preprocessor(Some("minimax_append_think"), Some("minimax_m2"));
let mut request: NvCreateChatCompletionRequest = serde_json::from_str(REQUEST_JSON).unwrap();
let tools: Vec<dynamo_protocols::types::ChatCompletionTool> =
serde_json::from_value(serde_json::json!([{
"type": "function",
"function": {
"name": "get_weather",
"description": "Get the current weather for a location.",
"parameters": {
"type": "object",
"properties": {"location": {"type": "string"}},
"required": ["location"]
}
}
}]))
.unwrap();
request.inner.tools = Some(tools);
request.inner.tool_choice = Some(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
let bare_params = r#"{"location": "Paris", "unit": "celsius"}"#;
let input_chunks = vec![mock_content_chunk(bare_params), mock_final_chunk()];
let input_stream = stream::iter(input_chunks.into_iter().map(Annotated::from_data));
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let output_chunks: Vec<Annotated<NvCreateChatCompletionStreamResponse>> =
output_stream.collect().await;
let mut reasoning = String::new();
let mut content = String::new();
let mut merged_tool_calls: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
if let Some(tcs) = &choice.delta.tool_calls {
for tc in tcs {
merged_tool_calls
.entry(tc.index)
.or_default()
.merge_from(tc);
}
}
}
}
assert!(
reasoning.is_empty(),
"reasoning_content must be empty (parser must be gated off), got: {reasoning:?}"
);
assert!(
!content.contains("<think>"),
"no <think> prefix should reach the client, got: {content:?}"
);
let tool_calls: Vec<MergedToolCall> = merged_tool_calls.values().cloned().collect();
assert_eq!(tool_calls.len(), 1, "expected one tool call");
assert_eq!(tool_calls[0].name.as_deref(), Some("get_weather"));
let args: Value = serde_json::from_str(&tool_calls[0].arguments).unwrap();
assert_eq!(
args,
serde_json::json!({"location": "Paris", "unit": "celsius"})
);
}
fn single_weather_tool() -> Vec<ChatCompletionTool> {
serde_json::from_value(serde_json::json!([{
"type": "function",
"function": {
"name": "get_weather",
"description": "Get the current weather for a location.",
"parameters": {
"type": "object",
"properties": {
"location": {"type": "string", "description": "City name"}
},
"required": ["location"]
}
}
}]))
.unwrap()
}
fn streaming_tool_request(
tool_choice: ChatCompletionToolChoiceOption,
) -> NvCreateChatCompletionRequest {
let mut request: NvCreateChatCompletionRequest = serde_json::from_value(serde_json::json!({
"model": "test-model",
"messages": [{"role": "user", "content": "What's the weather in San Francisco?"}],
"stream": true,
"temperature": 0.0
}))
.unwrap();
request.inner.tools = Some(single_weather_tool());
request.inner.tool_choice = Some(tool_choice);
request
}
fn streaming_json_schema_request(enable_thinking: bool) -> NvCreateChatCompletionRequest {
serde_json::from_value(serde_json::json!({
"model": "test-model",
"messages": [{"role": "user", "content": "Return a country and its capital."}],
"stream": true,
"temperature": 0.0,
"chat_template_kwargs": {"enable_thinking": enable_thinking},
"response_format": {
"type": "json_schema",
"json_schema": {
"name": "capital",
"schema": {
"type": "object",
"properties": {
"country": {"type": "string"},
"capital": {"type": "string"}
},
"required": ["country", "capital"],
"additionalProperties": false
}
}
}
}))
.unwrap()
}
fn streaming_json_object_request(enable_thinking: bool) -> NvCreateChatCompletionRequest {
serde_json::from_value(serde_json::json!({
"model": "test-model",
"messages": [{"role": "user", "content": "Return a JSON object."}],
"stream": true,
"temperature": 0.0,
"chat_template_kwargs": {"enable_thinking": enable_thinking},
"response_format": {
"type": "json_object"
}
}))
.unwrap()
}
fn enable_opt_in_reasoning(request: &mut NvCreateChatCompletionRequest, reasoning_parser: &str) {
if matches!(reasoning_parser, "deepseek_v3" | "deepseek_v3_1") {
request.chat_template_args =
Some(serde_json::from_value(serde_json::json!({"thinking": true})).unwrap());
} else if reasoning_parser == "mistral" {
request.chat_template_args =
Some(serde_json::from_value(serde_json::json!({"reasoning_effort": "high"})).unwrap());
}
}
struct DrainOutput {
reasoning: String,
content: String,
tool_calls: Vec<MergedToolCall>,
finish_reasons: Vec<FinishReason>,
}
async fn drain_stream(
output_stream: impl futures::Stream<Item = Annotated<NvCreateChatCompletionStreamResponse>>,
) -> DrainOutput {
let output_chunks: Vec<_> = Box::pin(output_stream).collect().await;
let mut reasoning = String::new();
let mut content = String::new();
let mut merged: BTreeMap<u32, MergedToolCall> = BTreeMap::new();
let mut finish_reasons = Vec::new();
for output in &output_chunks {
let Some(data) = output.data.as_ref() else {
continue;
};
for choice in &data.inner.choices {
if let Some(r) = &choice.delta.reasoning_content {
reasoning.push_str(r);
}
if let Some(c) = &choice.delta.content {
content.push_str(get_text(c));
}
if let Some(tcs) = &choice.delta.tool_calls {
for tc in tcs {
merged.entry(tc.index).or_default().merge_from(tc);
}
}
if let Some(fr) = choice.finish_reason {
finish_reasons.push(fr);
}
}
}
DrainOutput {
reasoning,
content,
tool_calls: merged.values().cloned().collect(),
finish_reasons,
}
}
fn assert_clean_tool_call(
case: &str,
content: &str,
tool_calls: &[MergedToolCall],
expected_location: &str,
) {
assert!(
!content.contains("get_weather"),
"{case}: tool-call JSON must not leak into content, got: {content:?}"
);
assert!(
!content.contains("<think>") && !content.contains("</think>"),
"{case}: think markers must not leak into content, got: {content:?}"
);
assert_eq!(tool_calls.len(), 1, "{case}: expected one tool call");
assert_eq!(
tool_calls[0].name.as_deref(),
Some("get_weather"),
"{case}: wrong tool name"
);
let args: Value = serde_json::from_str(&tool_calls[0].arguments)
.unwrap_or_else(|e| panic!("{case}: arguments not valid JSON: {e}"));
assert_eq!(
args,
serde_json::json!({"location": expected_location}),
"{case}: wrong arguments"
);
}
#[tokio::test]
async fn tool_choice_matrix_force_reasoning_required_bare_json() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
for &reasoning_parser in FORCE_REASONING_PARSERS {
for (case, prompt_injected) in [
(
"1a: force-reasoning + required + prompt_injected=false",
false,
),
(
"1b: force-reasoning + required + prompt_injected=true",
true,
),
] {
let case = format!("{case} + {reasoning_parser}");
let preprocessor = build_preprocessor(Some(reasoning_parser), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
enable_opt_in_reasoning(&mut request, reasoning_parser);
let input_stream = stream::iter(
vec![
mock_content_chunk(" \n"),
mock_content_chunk("["),
mock_content_chunk(&bare_json[1..]),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, prompt_injected, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not become reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(&case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
}
#[tokio::test]
async fn tool_choice_matrix_force_reasoning_named_bare_json() {
let bare_params = r#"{"location":"San Francisco"}"#;
for &reasoning_parser in FORCE_REASONING_PARSERS {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some("nemotron_nano"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Named(
ChatCompletionNamedToolChoice {
r#type: ChatCompletionToolType::Function,
function: FunctionName {
name: "get_weather".to_string(),
},
},
));
let input_stream = stream::iter(
vec![mock_content_chunk(bare_params), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = format!("2: force-reasoning + named + bare params + {reasoning_parser}");
assert!(
reasoning.is_empty(),
"{case}: reasoning_content must be empty, got: {reasoning:?}"
);
assert_clean_tool_call(&case, &content, &tool_calls, "San Francisco");
}
}
#[tokio::test]
async fn tool_choice_nemotron_v3_required_thinking_disabled_keeps_bare_json() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("nemotron_v3"), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
request.chat_template_args =
Some(serde_json::from_value(serde_json::json!({"enable_thinking": false})).unwrap());
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = "Nemotron v3 required + thinking disabled + bare JSON";
assert!(
reasoning.is_empty(),
"{case}: reasoning_content must be empty, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
}
#[tokio::test]
async fn tool_choice_deepseek_v3_required_thinking_disabled_keeps_bare_json() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
for reasoning_parser in ["deepseek_v3", "deepseek_v3_1", "deepseek_v3_2"] {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
request.chat_template_args =
Some(serde_json::from_value(serde_json::json!({"thinking": false})).unwrap());
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = format!("{reasoning_parser} required + thinking disabled + bare JSON");
assert!(
reasoning.is_empty(),
"{case}: reasoning_content must be empty, got: {reasoning:?}"
);
assert_clean_tool_call(&case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
#[tokio::test]
async fn tool_choice_force_reasoning_required_keeps_reasoning_before_guided_json() {
for &(reasoning_parser, close_marker) in REASONING_BEFORE_GUIDED_JSON_PARSERS {
for prompt_injected in [false, true] {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
enable_opt_in_reasoning(&mut request, reasoning_parser);
let split_at = close_marker.len() - 2;
let (close_prefix, close_suffix) = close_marker.split_at(split_at);
let close_and_json = format!(
r#"{close_suffix}[{{"name":"get_weather","parameters":{{"location":"San Francisco"}}}}]"#
);
let input_stream = stream::iter(
vec![
mock_content_chunk("Let me "),
mock_content_chunk("check."),
mock_content_chunk(close_prefix),
mock_content_chunk(&close_and_json),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, prompt_injected, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = format!(
"{reasoning_parser} required + reasoning boundary + prompt_injected={prompt_injected}"
);
assert_eq!(
reasoning, "Let me check.",
"{case}: reasoning_content must preserve the pre-boundary text"
);
assert_clean_tool_call(&case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
}
#[tokio::test]
async fn tool_choice_mistral_required_recognizes_split_reasoning_start() {
let preprocessor = build_preprocessor(Some("mistral"), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
enable_opt_in_reasoning(&mut request, "mistral");
let input_stream = stream::iter(
vec![
mock_content_chunk(" \n"),
mock_content_chunk("["),
mock_content_chunk("TH"),
mock_content_chunk("INK]Let me check.[/TH"),
mock_content_chunk(
r#"INK][{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#,
),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "Mistral required + whitespace + split [THINK] opener + guided JSON";
assert_eq!(reasoning, "Let me check.", "{case}: wrong reasoning");
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_force_reasoning_named_keeps_reasoning_before_guided_params() {
for &(reasoning_parser, close_marker) in REASONING_BEFORE_GUIDED_JSON_PARSERS {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some("nemotron_nano"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
enable_opt_in_reasoning(&mut request, reasoning_parser);
let input_stream = stream::iter(
vec![
mock_content_chunk("Let me check."),
mock_content_chunk(close_marker),
mock_content_chunk(r#"{"location":"San Francisco"}"#),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = format!("{reasoning_parser} named + reasoning boundary + guided params");
assert_eq!(
reasoning, "Let me check.",
"{case}: reasoning_content must preserve the pre-boundary text"
);
assert_clean_tool_call(&case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
#[tokio::test]
async fn tool_choice_matrix_non_force_required_no_injection_bare_json() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("qwen3"), Some("hermes"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = "3: non-force + required + prompt_injected=false + bare JSON";
assert!(
reasoning.is_empty(),
"{case}: parser must not produce reasoning when no <think> seen, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
}
#[tokio::test]
async fn tool_choice_matrix_non_force_required_prompt_injected_with_close_marker() {
let stream_text = r#"Let me check.</think>[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("qwen3"), Some("hermes"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(stream_text), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = "4: non-force + required + prompt_injected=true + reasoning</think>JSON";
assert_eq!(
reasoning.trim(),
"Let me check.",
"{case}: reasoning_content should hold only the pre-</think> text, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
}
#[tokio::test]
async fn tool_choice_matrix_non_force_required_prompt_injected_bare_json_contract() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("qwen3"), Some("hermes"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = "5 (contract): non-force + required + prompt_injected=true + bare JSON";
assert!(
tool_calls.is_empty(),
"{case}: contract case currently extracts no tool_calls (backend bug shape), got: {tool_calls:?}"
);
assert!(
content.is_empty(),
"{case}: content must remain empty (no leak), got: {content:?}"
);
assert!(
reasoning.contains("get_weather"),
"{case}: parser pins the JSON in reasoning_content under the broken contract, got: {reasoning:?}"
);
}
#[tokio::test]
async fn response_format_qwen3_prompt_injected_reasoning_then_json_preserves_channels() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let stream_text = format!("France is a country in Europe.</think>{json}");
let preprocessor = build_preprocessor(Some("qwen3"), None);
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(&stream_text), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert_eq!(reasoning.trim(), "France is a country in Europe.");
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_qwen3_prompt_injected_bare_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("qwen3"), None);
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"response_format JSON must not be reasoning_content, got: {reasoning:?}"
);
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_qwen3_no_thinking_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("qwen3"), None);
let request = streaming_json_schema_request(false);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(reasoning.is_empty());
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_gemma4_bare_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("gemma4"), None);
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"response_format JSON must not be reasoning_content, got: {reasoning:?}"
);
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_minimax_append_think_bare_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("minimax_append_think"), Some("minimax_m2"));
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"response_format JSON must not be reasoning_content, got: {reasoning:?}"
);
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_gpt_oss_bare_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("gpt_oss"), None);
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"response_format JSON must not be reasoning_content, got: {reasoning:?}"
);
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_gpt_oss_json_object_bare_json_stays_content() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let preprocessor = build_preprocessor(Some("gpt_oss"), None);
let request = streaming_json_object_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert!(
reasoning.is_empty(),
"response_format JSON must not be reasoning_content, got: {reasoning:?}"
);
assert_eq!(content, json);
}
#[tokio::test]
async fn response_format_gpt_oss_reasoning_then_json_preserves_channels() {
let json = r#"{"country":"France","capital":"Paris"}"#;
let stream_text = format!(
"<|channel|>analysis<|message|>Need answer as JSON.<|end|><|start|>assistant<|channel|>final<|message|>{json}"
);
let preprocessor = build_preprocessor(Some("gpt_oss"), None);
let request = streaming_json_schema_request(true);
let input_stream = stream::iter(
vec![mock_content_chunk(&stream_text), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning, content, ..
} = drain_stream(output_stream).await;
assert_eq!(reasoning, "Need answer as JSON.");
assert_eq!(content, json);
}
#[tokio::test]
async fn tool_choice_deepseek_v4_required_prompt_injected_bare_json_recovers() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("deepseek_v4"), Some("deepseek_v4"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "DeepSeek V4 required + prompt_injected=true + bare JSON";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_minimax_m3_required_prompt_injected_bare_json_recovers() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("minimax_m3"), Some("minimax_m3"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "MiniMax M3 required + prompt_injected=true + bare JSON";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_minimax_m2_required_keeps_reasoning_before_tool_xml() {
let preprocessor = build_preprocessor(Some("minimax_m2"), Some("minimax_m2"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let tool_call = "<minimax:tool_call>\
<invoke name=\"get_weather\"><parameter name=\"location\">San Francisco</parameter></invoke>\
</minimax:tool_call>";
let input_stream = stream::iter(
vec![
mock_content_chunk("I should call weather."),
mock_content_chunk("</think>"),
mock_content_chunk(tool_call),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "MiniMax M2 required + reasoning boundary + XML tool call";
assert_eq!(reasoning, "I should call weather.");
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(finish_reasons.contains(&FinishReason::ToolCalls));
}
#[tokio::test]
async fn tool_choice_minimax_m2_required_bare_json_bypasses_reasoning() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("minimax_m2"), Some("minimax_m2"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "MiniMax M2 required + bare guided JSON";
assert!(reasoning.is_empty());
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(finish_reasons.contains(&FinishReason::ToolCalls));
}
#[tokio::test]
async fn tool_choice_minimax_m2_required_thinking_disabled_keeps_tool_xml() {
let preprocessor = build_preprocessor(Some("minimax_m2"), Some("minimax_m2"));
let mut request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
request.chat_template_args =
Some(serde_json::from_value(serde_json::json!({"thinking": false})).unwrap());
let tool_call = "<minimax:tool_call>\
<invoke name=\"get_weather\"><parameter name=\"location\">San Francisco</parameter></invoke>\
</minimax:tool_call>";
let input_stream = stream::iter(
vec![mock_content_chunk(tool_call), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "MiniMax M2 required + thinking=false + XML tool call";
assert!(reasoning.is_empty());
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(finish_reasons.contains(&FinishReason::ToolCalls));
}
#[tokio::test]
async fn tool_calls_qwen3_coder_auto_routes_through_experimental_gate() {
let xml = "<tool_call>\n<function=get_weather>\n<parameter=location>\nSan Francisco\n</parameter>\n</function>\n</tool_call>";
let preprocessor = build_preprocessor(None, Some("qwen3_coder"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Auto);
let input_stream = stream::iter(
vec![mock_content_chunk(xml), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
content,
tool_calls,
finish_reasons,
..
} = drain_stream(output_stream).await;
let path = if std::env::var("DYN_ENABLE_EXPERIMENTAL_PARSERS_V2")
.is_ok_and(|v| matches!(v.trim(), "1" | "true" | "yes" | "on"))
{
"qwen3_coder auto -> dynamo-parsers-v2 (DYN_ENABLE_EXPERIMENTAL_PARSERS_V2 on)"
} else {
"qwen3_coder auto -> v1 jail (flag off)"
};
assert_clean_tool_call(path, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{path}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_prompt_injected_close_marker_json_keeps_reasoning_parser_for_dsv4_glm() {
let stream_text = r#"Let me check.</think>[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
for (case, reasoning_parser, tool_call_parser) in [
("DeepSeek V4", "deepseek_v4", "deepseek_v4"),
("GLM45", "glm45", "glm47"),
] {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some(tool_call_parser));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(stream_text), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
assert_eq!(
reasoning.trim(),
"Let me check.",
"{case}: reasoning_content should hold only the pre-</think> text, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
#[tokio::test]
async fn tool_choice_deepseek_v4_named_prompt_injected_bare_params_recovers() {
let bare_params = r#"{"location":"San Francisco"}"#;
let preprocessor = build_preprocessor(Some("deepseek_v4"), Some("deepseek_v4"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
let input_stream = stream::iter(
vec![mock_content_chunk(bare_params), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "DeepSeek V4 named + prompt_injected=true + bare params";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_minimax_m3_named_prompt_injected_bare_params_recovers() {
let bare_params = r#"{"location":"San Francisco"}"#;
let preprocessor = build_preprocessor(Some("minimax-m3"), Some("minimax-m3"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
let input_stream = stream::iter(
vec![mock_content_chunk(bare_params), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "MiniMax M3 named + prompt_injected=true + bare params";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_glm45_required_prompt_injected_bare_json_recovers() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("glm45"), Some("glm47"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![mock_content_chunk(bare_json), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "GLM45 required + prompt_injected=true + bare JSON";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_glm45_named_prompt_injected_bare_params_recovers() {
let bare_params = r#"{"location":"San Francisco"}"#;
let preprocessor = build_preprocessor(Some("glm45"), Some("glm47"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Named(
"get_weather".to_string().into(),
));
let input_stream = stream::iter(
vec![mock_content_chunk(bare_params), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, true, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
let case = "GLM45 named + prompt_injected=true + bare params";
assert!(
reasoning.is_empty(),
"{case}: guided JSON must not be classified as reasoning_content, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
#[tokio::test]
async fn tool_choice_structural_tag_keeps_reasoning_parser() {
for (case, reasoning_parser, tool_call_parser, structural_tool_call, prompt_injected) in [
(
"DeepSeek V4 DSML",
"deepseek_v4",
"deepseek_v4",
"<|DSML|tool_calls>\n\
<|DSML|invoke name=\"get_weather\">\n\
<|DSML|parameter name=\"location\" string=\"true\">San Francisco</|DSML|parameter>\n\
</|DSML|invoke>\n\
</|DSML|tool_calls>",
true,
),
(
"GLM XML",
"glm45",
"glm47",
"<tool_call>get_weather\
<arg_key>location</arg_key><arg_value>San Francisco</arg_value>\
</tool_call>",
true,
),
(
"Nemotron v3 Qwen3-Coder XML",
"nemotron_v3",
"qwen3_coder",
"<tool_call>\n\
<function=get_weather>\n\
<parameter=location>\n\
San Francisco\n\
</parameter>\n\
</function>\n\
</tool_call>",
false,
),
] {
let preprocessor = build_preprocessor(Some(reasoning_parser), Some(tool_call_parser));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let stream_text = format!("Let me check.</think>{structural_tool_call}");
let input_stream = stream::iter(
vec![mock_content_chunk(&stream_text), mock_final_chunk()]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, prompt_injected, true)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
finish_reasons,
} = drain_stream(output_stream).await;
assert_eq!(
reasoning.trim(),
"Let me check.",
"{case}: reasoning_content should hold only the pre-</think> text, got: {reasoning:?}"
);
assert!(
content.is_empty(),
"{case}: reasoning prefix or structural tags leaked into content: {content:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
assert!(
finish_reasons.contains(&FinishReason::ToolCalls),
"{case}: expected ToolCalls finish_reason, got: {finish_reasons:?}"
);
}
}
#[tokio::test]
async fn tool_choice_matrix_immediate_jail_reasoning_only_first_chunk() {
let bare_json = r#"[{"name":"get_weather","parameters":{"location":"San Francisco"}}]"#;
let preprocessor = build_preprocessor(Some("qwen3"), Some("hermes"));
let request = streaming_tool_request(ChatCompletionToolChoiceOption::Required);
let input_stream = stream::iter(
vec![
mock_reasoning_only_chunk("thinking briefly"),
mock_content_chunk(bare_json),
mock_final_chunk(),
]
.into_iter()
.map(Annotated::from_data),
);
let output_stream = preprocessor
.postprocessor_parsing_stream(input_stream, &request, false, false)
.expect("postprocessor_parsing_stream should build");
let DrainOutput {
reasoning,
content,
tool_calls,
..
} = drain_stream(output_stream).await;
let case = "6: Immediate jail + reasoning-only first chunk + JSON later";
assert!(
reasoning.contains("thinking briefly"),
"{case}: reasoning_content from the first chunk must reach the client, got: {reasoning:?}"
);
assert_clean_tool_call(case, &content, &tool_calls, "San Francisco");
}