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
101 changes: 49 additions & 52 deletions model_gateway/src/routers/grpc/regular/processor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,37 +97,35 @@ impl ResponseProcessor {
let mut processed_text = final_text;

if original_request.separate_reasoning && reasoning_parser_available {
let pooled_parser = utils::get_reasoning_parser(
// Fresh parser per request: non-streaming extraction keeps no state
// across requests, so avoid serializing on the shared pooled mutex.
if let Some(mut parser) = utils::create_reasoning_parser(
&self.reasoning_parser_factory,
self.configured_reasoning_parser.as_deref(),
&original_request.model,
);

let mut parser = pooled_parser.lock().await;
// Reset pooled parser to clean state before each request
parser.reset();

// If the template injected `<think>` in the prefill (thinking toggle
// is supported and effectively ON), start in reasoning mode.
if utils::should_mark_reasoning_started(
utils::extract_thinking_from_kwargs(
original_request.chat_template_kwargs.as_ref(),
tokenizer.as_ref(),
),
tokenizer.as_ref(),
) {
parser.mark_reasoning_started();
}
// If the template injected `<think>` in the prefill (thinking toggle
// is supported and effectively ON), start in reasoning mode.
if utils::should_mark_reasoning_started(
utils::extract_thinking_from_kwargs(
original_request.chat_template_kwargs.as_ref(),
tokenizer.as_ref(),
),
tokenizer.as_ref(),
) {
parser.mark_reasoning_started();
}

match parser.detect_and_parse_reasoning(&processed_text) {
Ok(result) => {
if !result.reasoning_text.is_empty() {
reasoning_text = Some(result.reasoning_text);
match parser.detect_and_parse_reasoning(&processed_text) {
Ok(result) => {
if !result.reasoning_text.is_empty() {
reasoning_text = Some(result.reasoning_text);
}
processed_text = result.normal_text;
}
Err(e) => {
warn!("Reasoning parsing error, skipping parsing: {e}");
}
processed_text = result.normal_text;
}
Err(e) => {
warn!("Reasoning parsing error, skipping parsing: {e}");
}
}
}
Expand Down Expand Up @@ -588,39 +586,38 @@ impl ResponseProcessor {
let mut processed_text = final_text;

if reasoning_parser_available {
let pooled_parser = utils::get_reasoning_parser(
// Fresh parser per request: non-streaming extraction keeps no state
// across requests, so avoid serializing on the shared pooled mutex.
if let Some(mut parser) = utils::create_reasoning_parser(
&self.reasoning_parser_factory,
self.configured_reasoning_parser.as_deref(),
&messages_request.model,
);
let mut parser = pooled_parser.lock().await;
// Reset pooled parser to clean state before each request
parser.reset();

// If thinking is effectively ON and template has a toggle, start in reasoning mode.
{
let user_thinking = match &messages_request.thinking {
Some(
messages::ThinkingConfig::Enabled { .. }
| messages::ThinkingConfig::Adaptive { .. },
) => Some(true),
Some(messages::ThinkingConfig::Disabled) => Some(false),
None => None,
};
if utils::should_mark_reasoning_started(user_thinking, tokenizer.as_ref()) {
parser.mark_reasoning_started();
) {
// If thinking is effectively ON and template has a toggle, start in reasoning mode.
{
let user_thinking = match &messages_request.thinking {
Some(
messages::ThinkingConfig::Enabled { .. }
| messages::ThinkingConfig::Adaptive { .. },
) => Some(true),
Some(messages::ThinkingConfig::Disabled) => Some(false),
None => None,
};
if utils::should_mark_reasoning_started(user_thinking, tokenizer.as_ref()) {
parser.mark_reasoning_started();
}
}
}

match parser.detect_and_parse_reasoning(&processed_text) {
Ok(result) => {
if !result.reasoning_text.is_empty() {
reasoning_text = Some(result.reasoning_text);
match parser.detect_and_parse_reasoning(&processed_text) {
Ok(result) => {
if !result.reasoning_text.is_empty() {
reasoning_text = Some(result.reasoning_text);
}
processed_text = result.normal_text;
}
Err(e) => {
warn!("Reasoning parsing error, skipping parsing: {e}");
}
processed_text = result.normal_text;
}
Err(e) => {
warn!("Reasoning parsing error, skipping parsing: {e}");
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion model_gateway/src/routers/grpc/utils/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,6 @@ pub(crate) use logprobs::{
pub(crate) use metrics::{error_type_from_status, route_to_endpoint};
pub(crate) use parsers::{
check_reasoning_parser_availability, check_tool_parser_availability, create_reasoning_parser,
create_tool_parser, extract_thinking_from_kwargs, get_reasoning_parser, get_tool_parser,
create_tool_parser, extract_thinking_from_kwargs, get_tool_parser,
should_mark_reasoning_started,
};
74 changes: 43 additions & 31 deletions model_gateway/src/routers/grpc/utils/parsers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,9 +4,7 @@ use llm_tokenizer::{
chat_template::{ThinkingKeyName, ThinkingToggle},
traits::Tokenizer,
};
use reasoning_parser::{
ParserFactory as ReasoningParserFactory, PooledParser as ReasoningPooledParser, ReasoningParser,
};
use reasoning_parser::{ParserFactory as ReasoningParserFactory, ReasoningParser};
use serde_json::Value;
use tool_parser::{
ParserFactory as ToolParserFactory, PooledParser as ToolPooledParser, ToolParser,
Expand Down Expand Up @@ -75,35 +73,10 @@ pub(crate) fn check_tool_parser_availability(
}
}

/// Get the appropriate reasoning parser for a model
/// Create a fresh reasoning parser instance.
///
/// If a parser name is explicitly configured, use that parser.
/// Otherwise, auto-detect based on the model name.
/// Get a pooled reasoning parser (for non-streaming where state doesn't matter)
pub(crate) fn get_reasoning_parser(
reasoning_parser_factory: &ReasoningParserFactory,
configured_parser: Option<&str>,
model: &str,
) -> ReasoningPooledParser {
if let Some(parser_name) = configured_parser {
// Use configured parser if specified
reasoning_parser_factory
.registry()
.get_pooled_parser(parser_name)
.unwrap_or_else(|| {
warn!(
"Configured reasoning parser '{}' not found, falling back to model-based selection",
parser_name
);
reasoning_parser_factory.get_pooled(model)
})
} else {
// Auto-detect based on model
reasoning_parser_factory.get_pooled(model)
}
}

/// Create a fresh reasoning parser instance (for streaming where state isolation is needed)
/// Used for both streaming (state isolation across chunks) and non-streaming
/// (avoids serializing on the shared pooled parser mutex).
pub(crate) fn create_reasoning_parser(
reasoning_parser_factory: &ReasoningParserFactory,
configured_parser: Option<&str>,
Expand Down Expand Up @@ -178,3 +151,42 @@ pub(crate) fn create_tool_parser(
tool_parser_factory.registry().create_for_model(model)
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn create_reasoning_parser_returns_independent_instances() {
let factory = ReasoningParserFactory::new();

// qwen3 starts with in_reasoning=false (explicit <think> required).
let mut a =
create_reasoning_parser(&factory, None, "qwen3").expect("qwen3 has a reasoning parser");
let mut b =
create_reasoning_parser(&factory, None, "qwen3").expect("qwen3 has a reasoning parser");

// Each call returns an independent instance: state mutated on one parser
// must not leak into the other (the shared pooled parser the non-streaming
// path used to take would have violated this).
a.mark_reasoning_started();
assert!(a.is_in_reasoning());
assert!(!b.is_in_reasoning());

// The untouched instance still parses a full document correctly.
let rb = b
.detect_and_parse_reasoning("<think>reasoning</think>answer")
.unwrap();
assert_eq!(rb.normal_text, "answer");
assert_eq!(rb.reasoning_text, "reasoning");
}

#[test]
fn create_reasoning_parser_honors_configured_parser() {
let factory = ReasoningParserFactory::new();

let parser = create_reasoning_parser(&factory, Some("qwen3"), "unknown-model")
.expect("configured qwen3 parser exists");
assert_eq!(parser.model_type(), "qwen3");
}
}
5 changes: 3 additions & 2 deletions model_gateway/src/routers/parse/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,14 +73,15 @@ pub async fn parse_reasoning(ctx: &Arc<AppContext>, req: &SeparateReasoningReque
);
};

let Some(pooled_parser) = factory.registry().get_pooled_parser(&req.reasoning_parser) else {
// Fresh parser per request: non-streaming extraction keeps no state across
// requests, so avoid serializing on the shared pooled mutex.
let Some(mut parser) = factory.registry().create_parser(&req.reasoning_parser) else {
return error_response(
StatusCode::BAD_REQUEST,
&format!("Unknown reasoning parser: {}", req.reasoning_parser),
);
};

let mut parser = pooled_parser.lock().await;
match parser.detect_and_parse_reasoning(&req.text) {
Ok(result) => (
StatusCode::OK,
Expand Down
Loading