Repository navigation
Connorli/fix func call parsing - #1045
ConnorLi96 wants to merge 3 commits into
Conversation
Signed-off-by: Scott Lee <scott@together.ai>
Signed-off-by: Scott Lee <scott@together.ai>
…ol-call tokens - Prioritize explicitly configured tool parser over JSON schema parsing - Support alternative delimiters (<|func_start|>/<|func_end|>) in KimiK2 parser - Prevent reasoning parser from consuming tool call markers - Strip leaked chatml tokens only when a parser is explicitly configured - Truncate trailing content after first valid JSON object for JSON response formats Signed-off-by: ConnorLi96 <ConnorLi96@users.noreply.github.com> Made-with: Cursor
📝 WalkthroughWalkthroughThis pull request introduces Prometheus multiprocess metrics collection for gRPC workers, enhances tool-call parsing with alternative delimiter support and improved reasoning detection, and updates gRPC router logic to prioritize configured parsers while normalizing responses and sanitizing ChatML tokens. Changes
Sequence Diagram(s)sequenceDiagram
participant ServeOrchestrator
participant gRPCWorker as gRPC Worker Process
participant WorkerManager
participant PythonCollector as Python Subprocess
ServeOrchestrator->>ServeOrchestrator: Create temp directory
ServeOrchestrator->>ServeOrchestrator: Set PROMETHEUS_MULTIPROC_DIR env var
ServeOrchestrator->>gRPCWorker: Launch with env var configured
gRPCWorker->>gRPCWorker: Write metrics to .db files<br/>in PROMETHEUS_MULTIPROC_DIR
WorkerManager->>WorkerManager: Detect gRPC workers present
WorkerManager->>PythonCollector: Spawn subprocess with<br/>prometheus_client collector
PythonCollector->>PythonCollector: Read PROMETHEUS_MULTIPROC_DIR
PythonCollector->>PythonCollector: Aggregate .db files into<br/>Prometheus text format
PythonCollector-->>WorkerManager: Return aggregated metrics
WorkerManager-->>WorkerManager: Append gRPC metrics<br/>to response
ServeOrchestrator->>ServeOrchestrator: On cleanup: Remove temp<br/>directory and .db files
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 inconclusive)
✅ Passed checks (2 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Code Review
This pull request implements gRPC worker metrics collection using Prometheus multiprocess directories and enhances tool and reasoning parsing, particularly for Kimi models. Key changes include a Python-based metrics aggregator in the worker manager, updated regex patterns for tool call detection, and logic to strip leaked ChatML tokens from model outputs. Review feedback highlights several critical improvements: the regex for function arguments needs to be more robust for nested JSON, and stripping special tokens via simple string replacement on streaming chunks is unreliable as tokens may be split across packets. Additionally, spawning a Python subprocess for metrics is inefficient, model-specific markers are hardcoded in base parsers, and token-stripping logic is duplicated across the codebase.
| // Pattern for complete tool calls | ||
| let tool_call_pattern = r"<\|tool_call_begin\|>\s*(?P<tool_call_id>[\w\.]+:\d+)\s*<\|tool_call_argument_begin\|>\s*(?P<function_arguments>\{.*?\})\s*<\|tool_call_end\|>"; | ||
| // Supports alternative delimiters: <|func_start|>/<|func_end|>; (?s) for multi-line JSON | ||
| let tool_call_pattern = r"(?s)<\|tool_call_begin\|>\s*(?P<tool_call_id>[\w\.]+:\d+)\s*(?:<\|tool_call_argument_begin\|>\s*|<\|func_start\|>\s*)?(?P<function_arguments>\{.*?\})\s*(?:<\|tool_call_end\|>|<\|func_end\|>)"; |
There was a problem hiding this comment.
The regex pattern \{.*?\} for function_arguments will fail to correctly capture nested JSON structures because the non-greedy match stops at the first encountered closing brace }. For streaming data where the buffer is cleared after each complete item, a greedy regex (.*) should be used to capture all accumulated content for the arguments until the end marker.
| let tool_call_pattern = r"(?s)<\|tool_call_begin\|>\s*(?P<tool_call_id>[\w\.]+:\d+)\s*(?:<\|tool_call_argument_begin\|>\s*|<\|func_start\|>\s*)?(?P<function_arguments>\{.*?\})\s*(?:<\|tool_call_end\|>|<\|func_end\|>)"; | |
| let tool_call_pattern = r"(?s)<\|tool_call_begin\|>\s*(?P<tool_call_id>[\w\.]+:\d+)\s*(?:<\|tool_call_argument_begin\|>\s*|<\|func_start\|>\s*)?(?P<function_arguments>\{.*)\s*(?:<\|tool_call_end\|>|<\|func_end\|>)"; |
References
- When parsing streaming data where the buffer is cleared after each complete item, a greedy regex (.*) can be intentionally used to capture all accumulated content for the current item's arguments.
| let mut delta = delta; | ||
| if self.configured_tool_parser.is_some() | ||
| || self.configured_reasoning_parser.is_some() | ||
| { | ||
| for token in [ | ||
| "<|im_end|>", "<|im_start|>", "<|im_user|>", | ||
| "<|im_assistant|>", "<|im_system|>", "<|im_middle|>", | ||
| "</think>", | ||
| ] { | ||
| delta = delta.replace(token, ""); | ||
| } | ||
| } |
There was a problem hiding this comment.
Using String::replace on individual streaming chunks is unreliable for stripping special tokens like <|im_end|> as they may be split across chunks. Furthermore, avoid inefficient string manipulations like buffer = buffer[pos + 1..].to_string() which reallocates memory. Instead, use a stateful buffer and buffer.drain(..pos + 1) to remove processed parts without reallocating the rest of the buffer.
References
- When processing streaming data, avoid inefficient string manipulations like buffer = buffer[pos + 1..].to_string() within a loop, as it reallocates memory. Instead, use more performant methods like buffer.drain(..pos + 1).
| if let Some(tool_pos) = processed_text.find("<|tool_calls_section_begin|>") { | ||
| let reasoning_text = processed_text[..tool_pos].trim().to_string(); | ||
| let normal_text = processed_text[tool_pos..].to_string(); | ||
| return Ok(ParserResult::new(normal_text, reasoning_text)); | ||
| } |
There was a problem hiding this comment.
The tool call section marker <|tool_calls_section_begin|> is hardcoded in the BaseReasoningParser. This marker is specific to certain models (like Kimi) and its presence in the base parser violates modularity. It should ideally be part of the ParserConfig so it can be customized per model without modifying the base implementation.
| async fn collect_prometheus_multiproc_metrics() -> Result<String, String> { | ||
| let dir = std::env::var("PROMETHEUS_MULTIPROC_DIR").map_err(|_| { | ||
| "PROMETHEUS_MULTIPROC_DIR not set; cannot collect metrics from gRPC workers".to_string() | ||
| })?; | ||
|
|
||
| let output = tokio::process::Command::new("python3") | ||
| .args([ | ||
| "-c", | ||
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", | ||
| ]) | ||
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | ||
| .output() | ||
| .await | ||
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | ||
|
|
||
| if !output.status.success() { | ||
| let stderr = String::from_utf8_lossy(&output.stderr); | ||
| return Err(format!("python3 prometheus collector failed: {stderr}")); | ||
| } | ||
|
|
||
| String::from_utf8(output.stdout) | ||
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}")) | ||
| } |
There was a problem hiding this comment.
Spawning a python3 subprocess for every metrics scrape is inefficient and introduces a heavy runtime dependency. Consider performing this aggregation natively in Rust. If parsing is required, prefer using a dedicated library like prometheus-parse over manual string manipulation or external scripts to ensure robustness.
References
- When parsing Prometheus metrics, prefer using a dedicated library (e.g., prometheus-parse) over fragile string manipulation methods to ensure robustness.
| if self.configured_tool_parser.is_some() || self.configured_reasoning_parser.is_some() { | ||
| for token in [ | ||
| "<|im_end|>", | ||
| "<|im_start|>", | ||
| "<|im_user|>", | ||
| "<|im_assistant|>", | ||
| "<|im_system|>", | ||
| "<|im_middle|>", | ||
| ] { | ||
| processed_text = processed_text.replace(token, ""); | ||
| } | ||
| processed_text = processed_text.trim().to_string(); | ||
| } |
There was a problem hiding this comment.
The logic for stripping leaked ChatML tokens is duplicated across multiple files and functions. This duplication makes the code harder to maintain. Extract this duplicated logic into a shared helper function to improve maintainability and reduce redundancy.
References
- Extract duplicated logic into a shared helper function to improve maintainability and reduce redundancy.
| for token in [ | ||
| "<|im_end|>", "<|im_start|>", "<|im_user|>", | ||
| "<|im_assistant|>", "<|im_system|>", "<|im_middle|>", | ||
| ] { |
There was a problem hiding this comment.
🟡 Nit: This token list is missing "</think>" which IS included in the other streaming stripping site (line 432). If a model emits </think> inside text that flows through the tool parser's normal_text output, it would leak to the client from this path but not the other.
Consider keeping the token lists consistent across all 4 stripping sites (2 in processor.rs, 2 in streaming.rs) — or better, extract a shared constant/helper to avoid drift.
| if self.configured_tool_parser.is_some() || self.configured_reasoning_parser.is_some() { | ||
| for token in [ | ||
| "<|im_end|>", | ||
| "<|im_start|>", | ||
| "<|im_user|>", | ||
| "<|im_assistant|>", | ||
| "<|im_system|>", | ||
| "<|im_middle|>", | ||
| ] { | ||
| processed_text = processed_text.replace(token, ""); | ||
| } | ||
| processed_text = processed_text.trim().to_string(); |
There was a problem hiding this comment.
🟡 Nit: This chatml-stripping block is now duplicated in 4 places (here, line ~720 in this file, and two sites in streaming.rs) with subtly different token lists — the streaming "regular content emission" path also strips </think> but these don't. Consider extracting a shared strip_leaked_tokens() helper with a single canonical token list to prevent drift between the sites.
| /// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess. | ||
| async fn collect_prometheus_multiproc_metrics() -> Result<String, String> { | ||
| let dir = std::env::var("PROMETHEUS_MULTIPROC_DIR").map_err(|_| { | ||
| "PROMETHEUS_MULTIPROC_DIR not set; cannot collect metrics from gRPC workers".to_string() | ||
| })?; | ||
|
|
||
| let output = tokio::process::Command::new("python3") | ||
| .args([ | ||
| "-c", | ||
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", | ||
| ]) | ||
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | ||
| .output() | ||
| .await | ||
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | ||
|
|
||
| if !output.status.success() { | ||
| let stderr = String::from_utf8_lossy(&output.stderr); | ||
| return Err(format!("python3 prometheus collector failed: {stderr}")); | ||
| } | ||
|
|
||
| String::from_utf8(output.stdout) | ||
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}")) | ||
| } | ||
|
|
There was a problem hiding this comment.
🟡 Nit: This spawns a python3 subprocess on every /metrics scrape. If Prometheus is polling at a typical 15s interval, that's 4 process spawns per minute (each importing prometheus_client). Probably fine for now, but worth keeping an eye on latency. If it becomes a bottleneck, consider caching the output for a short TTL or using a long-running sidecar.
| if is_json_response { | ||
| if processed_text.starts_with("```json") || processed_text.starts_with("```JSON") { | ||
| if let Some(start) = processed_text.find('\n') { | ||
| let inner = &processed_text[start + 1..]; | ||
| if let Some(end) = inner.rfind("```") { | ||
| processed_text = inner[..end].trim().to_string(); | ||
| } | ||
| } | ||
| } | ||
|
|
||
| if processed_text.starts_with('{') { | ||
| let mut depth = 0i32; | ||
| let mut in_string = false; | ||
| let mut escape = false; | ||
| let mut json_end = None; | ||
| for (i, ch) in processed_text.char_indices() { | ||
| if escape { escape = false; continue; } | ||
| if ch == '\\' && in_string { escape = true; continue; } | ||
| if ch == '"' { in_string = !in_string; continue; } | ||
| if in_string { continue; } | ||
| if ch == '{' { depth += 1; } else if ch == '}' { | ||
| depth -= 1; | ||
| if depth == 0 { json_end = Some(i + 1); break; } | ||
| } | ||
| } | ||
| if let Some(end) = json_end { | ||
| processed_text = processed_text[..end].to_string(); | ||
| } | ||
| } |
There was a problem hiding this comment.
🟡 Nit: Two minor edge cases in the JSON cleanup:
rfind("```")(line 219) could match triple-backticks inside a JSON string value, truncating the response. Unlikely but possible.- The depth-counting trimmer only handles
{roots — a JSON Schema response that is an array ([...]) would not be trimmed. Most schemas produce objects, butJsonObjectformat technically allows arrays.
Neither is likely to bite in practice, just flagging for awareness.
There was a problem hiding this comment.
Looks good! The changes are well-structured: tool parser priority fix, alternative delimiter support in KimiK2, chatml token stripping, JSON response cleanup, and gRPC prometheus metrics collection all make sense.
4 nits posted (no blockers):
- 🟡 x4: Inconsistent
</think>in one of the 4 chatml-stripping sites, suggest extracting a shared helper to prevent drift, python3 subprocess per metrics scrape worth monitoring, minor JSON cleanup edge cases.
0 🔴 Important · 4 🟡 Nit · 0 🟣 Pre-existing
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d8616a8648
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| let tool_chunks = if is_specific_function | ||
| && !(self.configured_tool_parser.is_some() | ||
| && tool_parser_available) | ||
| { |
There was a problem hiding this comment.
Preserve configured parser precedence in chat streaming
When ToolChoice::Function is used and a model-specific parser is configured, this branch now routes into process_tool_calls_stream, but used_json_schema remains true for that mode, so the downstream parser creation path still selects Some("json") instead of the configured native parser. In practice, models that emit native tool-call tokens (e.g. <|tool_call_begin|>...) will still fail to emit proper tool_call deltas in streaming chat completions, so the parser-priority fix is not actually applied for this case.
Useful? React with 👍 / 👎.
| if is_specific_function | ||
| && !(self.configured_tool_parser.is_some() && tool_parser_available) |
There was a problem hiding this comment.
Preserve configured parser precedence in messages streaming
This new guard also diverts ToolChoice::Tool requests to the incremental parser path when a configured parser exists, but the parser instance for this flow is still initialized as JSON whenever used_json_schema is true. That means configured native parsers are bypassed in the exact path this change introduced, so models that stream native tool-call tokens continue to miss/garble tool-use events in the Messages streaming API.
Useful? React with 👍 / 👎.
|
|
||
| // Pattern for removing completed tool calls | ||
| let end_pattern = r"<\|tool_call_begin\|>.*?<\|tool_call_end\|>"; | ||
| let end_pattern = r"<\|tool_call_begin\|>.*?(?:<\|tool_call_end\|>|<\|func_end\|>)"; |
There was a problem hiding this comment.
Make Kimi end-token regex handle multiline arguments
The stream parser now supports multiline JSON via (?s) in the start/capture regexes, but this end-pattern still does not use dotall mode. For multiline arguments, tool_call_end_pattern.find(...) can fail to match completed calls, and parse_incremental then clears the entire buffer, which drops any trailing content (including a following tool call) from the same chunk.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Actionable comments posted: 2
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (3)
model_gateway/src/routers/grpc/regular/streaming.rs (2)
380-408:⚠️ Potential issue | 🟠 MajorConfigured parsers still lose to JSON parsing in the streaming tool paths.
In chat streaming, this still passes
used_json_schema = trueintoprocess_tool_calls_stream, and Line 1240 inside that helper turns that intoSome("json"). In Messages streaming, falling through here does not help either because Lines 1665-1669 already builtstreaming_tool_parserasSome("json")whenused_json_schemawas true. Sotool_choice=function/required/anystill bypasses the configured parser in streaming.🛠️ Suggested chat-path fix
- used_json_schema, + used_json_schema + && !(self.configured_tool_parser.is_some() + && tool_parser_available),Apply the same guard where
streaming_tool_parseris constructed for the Messages path soToolChoice::Tool/Anyalso stop forcingSome("json").Also applies to: 1779-1835
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/grpc/regular/streaming.rs` around lines 380 - 408, The streaming paths are incorrectly forcing a JSON parser when used_json_schema == true, bypassing configured parsers; update the Messages and chat streaming branches so that when tool_choice is Tool/Any (ToolChoice::Tool / ToolChoice::Any) you do not coerce streaming_tool_parser or the used_json_schema flag into Some("json") — mirror the guard used in the chat-path so that process_tool_calls_stream is not passed used_json_schema=true (and streaming_tool_parser is not set to "json") when a configured parser exists; specifically adjust the construction of streaming_tool_parser and the call-sites of process_tool_calls_stream and process_specific_function_stream to honor configured_tool_parser and tool_parser_available instead of forcing JSON.
424-446:⚠️ Potential issue | 🟠 MajorDon't emit original logprobs after rewriting
delta.This block strips tokens out of
delta, but Lines 441-445 still attach the unfilteredchoice_logprobs. When the cleanup changes the chunk, the streamed logprobs refer to text the client never received.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/grpc/regular/streaming.rs` around lines 424 - 446, The code strips parser-only tokens from the chunk `delta` but still attaches the original `choice_logprobs` via `ChatCompletionStreamResponse::builder(...).add_choice_content_with_logprobs(...)`, causing logprobs to refer to text the client never received; update the streaming logic in the block that handles `delta` so that when `delta` was modified (i.e., when `self.configured_tool_parser` or `self.configured_reasoning_parser` is Some and any of the tokens were removed) you do not send the original `choice_logprobs` — instead either recompute/trim logprobs to match the filtered `delta` or call the variant that sends content without logprobs (e.g., use the `add_choice_content` path) so the client only receives logprobs that correspond to the emitted text.model_gateway/src/routers/grpc/regular/processor.rs (1)
191-244:⚠️ Potential issue | 🟠 MajorDrop or recompute logprobs after this output cleanup.
Line 191 builds
logprobsbefore this block strips ChatML tokens and truncates fenced/trailing JSON. If the cleanup changesprocessed_text,ChatChoice.logprobsno longer describes the returnedmessage.content.🛠️ Safe fallback
- let logprobs = complete.output_logprobs().map(|ref proto_logprobs| { + let mut logprobs = complete.output_logprobs().map(|ref proto_logprobs| { utils::convert_proto_to_openai_logprobs(proto_logprobs, tokenizer) }); + let pre_cleanup_text = processed_text.clone();+ if processed_text != pre_cleanup_text { + logprobs = None; + }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/grpc/regular/processor.rs` around lines 191 - 244, logprobs is computed from complete.output_logprobs() before you mutate processed_text, so ChatChoice.logprobs can become inconsistent; after the cleanup that strips ChatML tokens / truncates fenced/trailing JSON (the block referencing self.configured_tool_parser / self.configured_reasoning_parser and the JSON-trimming logic), either recompute or drop the logprobs: keep a copy of the proto logprobs (complete.output_logprobs()) if you can recompute to match the final processed_text, otherwise detect if processed_text was modified and set the final ChatChoice.logprobs to None (or omit it) instead of using the stale logprobs variable so ChatChoice.logprobs always matches the returned message.content.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@crates/tool_parser/src/parsers/kimik2.rs`:
- Around line 67-69: The regex stored in end_pattern used to build
tool_call_end_pattern lacks the (?s) DOTALL flag, so Regex::new(end_pattern) in
tool_call_end_pattern doesn't match multi-line tool call bodies; update
end_pattern to include the (?s) prefix (e.g.,
"(?s)<\\|tool_call_begin\\|>.*?(?:<\\|tool_call_end\\|>|<\\|func_end\\|>)") and
rebuild tool_call_end_pattern with Regex::new so tool_call_end_pattern.find(...)
will correctly match tool calls that span multiple lines.
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 88-116: Add a doc comment to the
collect_prometheus_multiproc_metrics function stating its runtime prerequisites:
that a python3 executable must be available on PATH, the Python
prometheus_client package (with multiprocess support) must be installed, and
PROMETHEUS_MULTIPROC_DIR must be set and pointing to the multiprocess metrics
directory; include a brief note on expected failure modes (missing python3,
missing package, or unset env) so callers know why this function can return an
Err and how to remediate.
---
Outside diff comments:
In `@model_gateway/src/routers/grpc/regular/processor.rs`:
- Around line 191-244: logprobs is computed from complete.output_logprobs()
before you mutate processed_text, so ChatChoice.logprobs can become
inconsistent; after the cleanup that strips ChatML tokens / truncates
fenced/trailing JSON (the block referencing self.configured_tool_parser /
self.configured_reasoning_parser and the JSON-trimming logic), either recompute
or drop the logprobs: keep a copy of the proto logprobs
(complete.output_logprobs()) if you can recompute to match the final
processed_text, otherwise detect if processed_text was modified and set the
final ChatChoice.logprobs to None (or omit it) instead of using the stale
logprobs variable so ChatChoice.logprobs always matches the returned
message.content.
In `@model_gateway/src/routers/grpc/regular/streaming.rs`:
- Around line 380-408: The streaming paths are incorrectly forcing a JSON parser
when used_json_schema == true, bypassing configured parsers; update the Messages
and chat streaming branches so that when tool_choice is Tool/Any
(ToolChoice::Tool / ToolChoice::Any) you do not coerce streaming_tool_parser or
the used_json_schema flag into Some("json") — mirror the guard used in the
chat-path so that process_tool_calls_stream is not passed used_json_schema=true
(and streaming_tool_parser is not set to "json") when a configured parser
exists; specifically adjust the construction of streaming_tool_parser and the
call-sites of process_tool_calls_stream and process_specific_function_stream to
honor configured_tool_parser and tool_parser_available instead of forcing JSON.
- Around line 424-446: The code strips parser-only tokens from the chunk `delta`
but still attaches the original `choice_logprobs` via
`ChatCompletionStreamResponse::builder(...).add_choice_content_with_logprobs(...)`,
causing logprobs to refer to text the client never received; update the
streaming logic in the block that handles `delta` so that when `delta` was
modified (i.e., when `self.configured_tool_parser` or
`self.configured_reasoning_parser` is Some and any of the tokens were removed)
you do not send the original `choice_logprobs` — instead either recompute/trim
logprobs to match the filtered `delta` or call the variant that sends content
without logprobs (e.g., use the `add_choice_content` path) so the client only
receives logprobs that correspond to the emitted text.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 712739d4-f253-4af5-9c4d-24495824a7ab
📒 Files selected for processing (6)
bindings/python/src/smg/serve.pycrates/reasoning_parser/src/parsers/base.rscrates/tool_parser/src/parsers/kimik2.rsmodel_gateway/src/core/worker_manager.rsmodel_gateway/src/routers/grpc/regular/processor.rsmodel_gateway/src/routers/grpc/regular/streaming.rs
| // Pattern for removing completed tool calls | ||
| let end_pattern = r"<\|tool_call_begin\|>.*?<\|tool_call_end\|>"; | ||
| let end_pattern = r"<\|tool_call_begin\|>.*?(?:<\|tool_call_end\|>|<\|func_end\|>)"; | ||
| let tool_call_end_pattern = Regex::new(end_pattern).expect("Valid regex pattern"); |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
# First, locate the kimik2.rs file
fd -t f 'kimik2\.rs'Repository: lightseekorg/smg
Length of output: 147
🏁 Script executed:
# Read the relevant sections of kimik2.rs
cat -n crates/tool_parser/src/parsers/kimik2.rs | sed -n '50,100p'Repository: lightseekorg/smg
Length of output: 2682
🏁 Script executed:
# Now read the buffer-drain logic around lines 280-300
cat -n crates/tool_parser/src/parsers/kimik2.rs | sed -n '280,310p'Repository: lightseekorg/smg
Length of output: 1879
🏁 Script executed:
# Check lines 260-275 to verify the completeness check
cat -n crates/tool_parser/src/parsers/kimik2.rs | sed -n '260,275p'Repository: lightseekorg/smg
Length of output: 1002
🏁 Script executed:
# Let me verify the proposed fix by checking if there are any other usages of tool_call_end_pattern
rg -n "tool_call_end_pattern" crates/tool_parser/src/parsers/kimik2.rsRepository: lightseekorg/smg
Length of output: 327
🏁 Script executed:
# Check test file to understand expected multi-line scenarios
cat -n crates/tool_parser/tests/tool_parser_kimik2.rs | head -100Repository: lightseekorg/smg
Length of output: 3994
🏁 Script executed:
# Look for tests with multi-line JSON arguments
rg -n "\\\\n" crates/tool_parser/tests/tool_parser_kimik2.rs | head -20Repository: lightseekorg/smg
Length of output: 264
🏁 Script executed:
# Search for tests with actual newlines inside JSON payloads
rg -B5 -A5 'tool_call_argument_begin.*?\{' crates/tool_parser/tests/tool_parser_kimik2.rs | head -50Repository: lightseekorg/smg
Length of output: 2092
🏁 Script executed:
# Verify the Regex behavior by checking if there are any multiline payload tests
rg -n "func_end" crates/tool_parser/tests/tool_parser_kimik2.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check if there are any tests using <|func_start|>/<|func_end|> delimiters
rg -n "func_start" crates/tool_parser/tests/tool_parser_kimik2.rsRepository: lightseekorg/smg
Length of output: 42
🏁 Script executed:
# Check the is_complete_json logic to see when buffer drain is triggered
cat -n crates/tool_parser/src/parsers/kimik2.rs | sed -n '270,295p'Repository: lightseekorg/smg
Length of output: 1722
Add (?s) flag to tool_call_end_pattern to match multi-line JSON arguments.
The extractor regexes at lines 61 and 64 use (?s) to accept multi-line arguments, but tool_call_end_pattern at line 68 does not. When a completed tool call spans multiple lines, tool_call_end_pattern.find() at line 292 fails to match (since . doesn't cross newlines), causing line 296 to clear the entire buffer and lose trailing content or subsequent tool calls.
🐛 Fix
- let end_pattern = r"<\|tool_call_begin\|>.*?(?:<\|tool_call_end\|>|<\|func_end\|>)";
+ let end_pattern = r"(?s)<\|tool_call_begin\|>.*?(?:<\|tool_call_end\|>|<\|func_end\|>)";🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@crates/tool_parser/src/parsers/kimik2.rs` around lines 67 - 69, The regex
stored in end_pattern used to build tool_call_end_pattern lacks the (?s) DOTALL
flag, so Regex::new(end_pattern) in tool_call_end_pattern doesn't match
multi-line tool call bodies; update end_pattern to include the (?s) prefix
(e.g.,
"(?s)<\\|tool_call_begin\\|>.*?(?:<\\|tool_call_end\\|>|<\\|func_end\\|>)") and
rebuild tool_call_end_pattern with Regex::new so tool_call_end_pattern.find(...)
will correctly match tool calls that span multiple lines.
| /// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess. | ||
| async fn collect_prometheus_multiproc_metrics() -> Result<String, String> { | ||
| let dir = std::env::var("PROMETHEUS_MULTIPROC_DIR").map_err(|_| { | ||
| "PROMETHEUS_MULTIPROC_DIR not set; cannot collect metrics from gRPC workers".to_string() | ||
| })?; | ||
|
|
||
| let output = tokio::process::Command::new("python3") | ||
| .args([ | ||
| "-c", | ||
| "import sys\n\ | ||
| from prometheus_client import CollectorRegistry, generate_latest\n\ | ||
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | ||
| registry = CollectorRegistry()\n\ | ||
| MultiProcessCollector(registry)\n\ | ||
| sys.stdout.buffer.write(generate_latest(registry))\n", | ||
| ]) | ||
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | ||
| .output() | ||
| .await | ||
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | ||
|
|
||
| if !output.status.success() { | ||
| let stderr = String::from_utf8_lossy(&output.stderr); | ||
| return Err(format!("python3 prometheus collector failed: {stderr}")); | ||
| } | ||
|
|
||
| String::from_utf8(output.stdout) | ||
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}")) | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Runtime dependency on Python requires documentation.
This function introduces a runtime dependency on python3 being available in PATH and the prometheus_client package being installed. While this aligns with the Python orchestrator in serve.py that sets up PROMETHEUS_MULTIPROC_DIR, the dependency should be documented.
Consider adding a doc comment noting the prerequisites:
📝 Suggested documentation
-/// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess.
+/// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess.
+///
+/// # Prerequisites
+/// - `PROMETHEUS_MULTIPROC_DIR` environment variable must be set (typically by the Python orchestrator)
+/// - `python3` must be available in PATH
+/// - `prometheus_client` Python package must be installed
+///
+/// # Errors
+/// Returns an error string if the env var is missing, subprocess fails, or output is not valid UTF-8.
async fn collect_prometheus_multiproc_metrics() -> Result<String, String> {📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| /// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess. | |
| async fn collect_prometheus_multiproc_metrics() -> Result<String, String> { | |
| let dir = std::env::var("PROMETHEUS_MULTIPROC_DIR").map_err(|_| { | |
| "PROMETHEUS_MULTIPROC_DIR not set; cannot collect metrics from gRPC workers".to_string() | |
| })?; | |
| let output = tokio::process::Command::new("python3") | |
| .args([ | |
| "-c", | |
| "import sys\n\ | |
| from prometheus_client import CollectorRegistry, generate_latest\n\ | |
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | |
| registry = CollectorRegistry()\n\ | |
| MultiProcessCollector(registry)\n\ | |
| sys.stdout.buffer.write(generate_latest(registry))\n", | |
| ]) | |
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | |
| .output() | |
| .await | |
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | |
| if !output.status.success() { | |
| let stderr = String::from_utf8_lossy(&output.stderr); | |
| return Err(format!("python3 prometheus collector failed: {stderr}")); | |
| } | |
| String::from_utf8(output.stdout) | |
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}")) | |
| } | |
| /// Collect gRPC worker metrics by aggregating `PROMETHEUS_MULTIPROC_DIR` via a python3 subprocess. | |
| /// | |
| /// # Prerequisites | |
| /// - `PROMETHEUS_MULTIPROC_DIR` environment variable must be set (typically by the Python orchestrator) | |
| /// - `python3` must be available in PATH | |
| /// - `prometheus_client` Python package must be installed | |
| /// | |
| /// # Errors | |
| /// Returns an error string if the env var is missing, subprocess fails, or output is not valid UTF-8. | |
| async fn collect_prometheus_multiproc_metrics() -> Result<String, String> { | |
| let dir = std::env::var("PROMETHEUS_MULTIPROC_DIR").map_err(|_| { | |
| "PROMETHEUS_MULTIPROC_DIR not set; cannot collect metrics from gRPC workers".to_string() | |
| })?; | |
| let output = tokio::process::Command::new("python3") | |
| .args([ | |
| "-c", | |
| "import sys\n\ | |
| from prometheus_client import CollectorRegistry, generate_latest\n\ | |
| from prometheus_client.multiprocess import MultiProcessCollector\n\ | |
| registry = CollectorRegistry()\n\ | |
| MultiProcessCollector(registry)\n\ | |
| sys.stdout.buffer.write(generate_latest(registry))\n", | |
| ]) | |
| .env("PROMETHEUS_MULTIPROC_DIR", &dir) | |
| .output() | |
| .await | |
| .map_err(|e| format!("failed to run python3 prometheus collector: {e}"))?; | |
| if !output.status.success() { | |
| let stderr = String::from_utf8_lossy(&output.stderr); | |
| return Err(format!("python3 prometheus collector failed: {stderr}")); | |
| } | |
| String::from_utf8(output.stdout) | |
| .map_err(|e| format!("prometheus collector output is not valid UTF-8: {e}")) | |
| } |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/worker_manager.rs` around lines 88 - 116, Add a doc
comment to the collect_prometheus_multiproc_metrics function stating its runtime
prerequisites: that a python3 executable must be available on PATH, the Python
prometheus_client package (with multiprocess support) must be installed, and
PROMETHEUS_MULTIPROC_DIR must be set and pointing to the multiprocess metrics
directory; include a brief note on expected failure modes (missing python3,
missing package, or unset env) so callers know why this function can return an
Err and how to remediate.
Description
Problem
Function calling is broken for models that emit native tool-call tokens (e.g. Kimi K2). Three issues compound:
tool_choiceisrequiredor a specific function, the code picks the JSON schema parser over the model-specific tool parser. Models like Kimi always output native tokens (<|tool_call_begin|>, etc.) regardless oftool_choice, so the JSON schema path fails silently.</think>and goes straight to<|tool_calls_section_begin|>, the reasoning parser treats the entire output (including tool-call tokens) as reasoning content. The downstream tool parser never sees them.<|func_start|>/<|func_end|>instead of<|tool_call_argument_begin|>/<|tool_call_end|>, and may produce multi-line JSON arguments. The KimiK2 parser regex rejects all of these.Solution
<|tool_calls_section_begin|>before classifying content as reasoning, in both streaming and non-streaming paths.<|func_start|>/<|func_end|>as alternative delimiters; add(?s)flag for multi-line JSON arguments.<|im_end|>, etc.) only when a model-specific parser is explicitly configured, to avoid affecting models that legitimately output these tokens.json_object/json_schemaresponse formats.Changes
crates/reasoning_parser/src/parsers/base.rs: bail out of reasoning mode when tool-call section markers are found (streaming + non-streaming)crates/tool_parser/src/parsers/kimik2.rs: support<|func_start|>/<|func_end|>delimiters,(?s)for multi-line JSON, consistent end-token cleanupmodel_gateway/src/routers/grpc/regular/processor.rs: three-tier parser priority, conditional chatml token stripping, JSON boundary truncation (Chat Completions + Messages API)model_gateway/src/routers/grpc/regular/streaming.rs: same priority fix for streaming paths, conditional chatml/think token stripping in content deltasTest Plan
tool_choice: "auto","required", and specific function — all correctly parse native tool-call tokenstool_choiceis setChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Bug Fixes