feat(gateway): add MessageResponseProcessingStage for Messages API (non-streaming) - #747
Conversation
…on-streaming) Add Stage 7 (response processing) for the Messages API gRPC pipeline. This converts backend ProtoGenerateComplete responses into Anthropic Message format with proper ContentBlock construction and StopReason mapping. Non-streaming only; streaming deferred to follow-up PR. What changed: - processor.rs: add process_non_streaming_messages_response() to ResponseProcessor — full pipeline: token decoding, reasoning parsing, tool call parsing, content block construction (Thinking → Text → ToolUse), StopReason mapping (EndTurn/MaxTokens/StopSequence/ToolUse), and messages::Usage building - messages/response_processing.rs: new MessageResponseProcessingStage that extracts execution result, dispatch metadata, tokenizer, and stop decoder from RequestContext, delegates to ResponseProcessor, and stores FinalResponse::Messages - message_utils.rs: add get_history_tool_calls_count_messages() for counting tool use blocks in Messages API request history (needed for KimiK2-style tool call ID generation) - messages/mod.rs: wire response_processing module with unused_imports expect (wired in pipeline factory PR) Why: This is the fourth PR in the Messages API gRPC support series. With preparation (PR #741), request building (PR #744), and now response processing, three of the four endpoint-specific pipeline stages are complete. The shared stages (worker selection, client acquisition, dispatch, execution) are reused from the existing pipeline. How: Follows the same architecture as chat's response processing but adapted for Anthropic Message types: - Reuses existing convert_message_tool_choice() from message_utils to bridge Messages ToolChoice → Chat ToolChoice for parse_json_schema_response - Reuses ResponseProcessor's parse_tool_calls() for model-predicted path - Content blocks ordered per Anthropic convention: Thinking first, Text, then ToolUse blocks - Tool calls parsed as OpenAI ToolCall (via existing parsers) then converted to ContentBlock::ToolUse with JSON input - Messages always n=1, no logprobs - ThinkingConfig::Enabled check replaces separate_reasoning bool Refs: #739, #741, #744 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request significantly advances the Messages API gRPC support by implementing the crucial response processing stage for non-streaming requests. It establishes the pipeline to transform raw model outputs into the structured Anthropic Message format, handling token decoding, parsing of reasoning and tool calls, and accurate mapping of stop reasons and usage statistics. This lays the groundwork for full Messages API functionality, with streaming support planned for a subsequent PR. Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for GitHub and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
📝 WalkthroughWalkthroughAdds a non-streaming Messages API response processor and pipeline stage that converts ExecutionResult into Anthropic-like Message objects, plus utilities to count historical ToolUse blocks in CreateMessageRequest inputs. Changes
Sequence DiagramsequenceDiagram
participant Ctx as RequestContext
participant Stage as MessageResponseProcessingStage
participant Processor as ResponseProcessor
participant Tokenizer as Tokenizer
participant Decoder as StopSequenceDecoder
participant Final as FinalResponse::Messages
Ctx->>Stage: execute(ctx)
Stage->>Ctx: read execution_result, dispatch
Stage->>Ctx: get cached Tokenizer
Stage->>Ctx: get StopSequenceDecoder
Stage->>Processor: process_non_streaming_messages_response(exec_result, request, dispatch, tokenizer, decoder)
Processor->>Tokenizer: decode/convert tokens
Processor->>Decoder: apply stop sequence decoding
Processor->>Processor: parse reasoning (if present)
Processor->>Processor: parse tool calls (JSON-schema or parser)
Processor->>Processor: assemble content blocks (Thinking, Text, ToolUse)
Processor-->>Stage: Message
Stage->>Ctx: store FinalResponse::Messages(Message)
Stage-->>Ctx: Ok(None)
Estimated Code Review Effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly Related PRs
Suggested Labels
Suggested Reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
📝 Coding Plan
Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces the response processing stage for the non-streaming Messages API, a crucial part of the gRPC support series. The implementation is well-structured, follows existing patterns, and effectively reuses components. The two identified areas for improvement, enhancing debuggability by logging a silent failure and refactoring a function for better clarity and idiomatic Rust, are valid and align with repository guidelines. Overall, this is a solid and well-executed feature addition.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 291052716e
ℹ️ 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".
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/routers/grpc/regular/processor.rs`:
- Around line 458-468: In process_non_streaming_messages_response, don't
silently take the first entry from all_responses; instead explicitly enforce the
n=1 invariant by checking all_responses length (or pattern) and
returning/logging an internal error if there are zero or more than one
responses; update the existing ok_or_else branch on all_responses (and the error
message currently using error::internal_error("no_responses", ...)) to also
detect multiple ProtoGenerateComplete entries and call error::internal_error
with a distinct code (e.g., "multiple_responses") and descriptive message so
extra backend completions are rejected rather than ignored.
- Around line 635-649: Compute matched_stop_json() once and derive both
stop_sequence and stop_reason from that single value so they remain consistent:
first call complete.matched_stop_json() into a local (e.g. matched_stop_json),
convert it to a string (matched_stop_str) only if it's a JSON string, set
stop_sequence = matched_stop_str.clone(), then compute stop_reason using the
same inputs in order—if tool_calls.is_some() => StopReason::ToolUse (and ensure
stop_sequence is None), else if matched_stop_str.is_some() =>
StopReason::StopSequence, else if finish_reason_str == "length" =>
StopReason::MaxTokens, else => StopReason::EndTurn—so StopSequence is only set
when stop_reason is StopSequence and non-string JSON does not produce a
stop_sequence.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 96bf46ea-e469-4628-8b93-848613f0705b
📒 Files selected for processing (4)
model_gateway/src/routers/grpc/regular/processor.rsmodel_gateway/src/routers/grpc/regular/stages/messages/mod.rsmodel_gateway/src/routers/grpc/regular/stages/messages/response_processing.rsmodel_gateway/src/routers/grpc/utils/message_utils.rs
processor.rs: - Enforce n=1 invariant: reject 0 or >1 responses with distinct errors instead of silently taking the first - Log warning on tool call argument JSON parse failure instead of silent fallback to empty object - Derive stop_reason and stop_sequence from same conditions to prevent inconsistency (e.g. ToolUse + stop_sequence, or StopSequence + None) - Map finish_reason="tool_calls" to StopReason::ToolUse even when tool parsing was skipped/failed message_utils.rs: - Simplify get_history_tool_calls_count_messages with flat_map instead of nested filter_map + count + sum Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/routers/grpc/regular/processor.rs`:
- Around line 486-487: Update the attribute reason on the guarded unwrap in
processor.rs to use the repo convention "INVARIANT:"; specifically, change the
#[expect(clippy::unwrap_used, reason = "...")] attached to the let complete =
all_responses.into_iter().next().unwrap(); line so the reason begins with
"INVARIANT:" (e.g., "INVARIANT: checked len == 1 above") to document the
safe-code assumption consistently.
- Around line 584-613: Compute history tool-call count once before the
tool-choice branch and reuse it in both parsing paths: call
utils::message_utils::get_history_tool_calls_count_messages(&messages_request)
into a local let (e.g., history_tool_calls_count) before evaluating
used_json_schema or tool_parser_available, then pass that variable to
utils::parse_json_schema_response(...) and to self.parse_tool_calls(...). Update
references in the block around used_json_schema, chat_tool_choice,
parse_json_schema_response, and parse_tool_calls to use the new local instead of
calling get_history_tool_calls_count_messages(&messages_request) twice.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 155affe6-9370-4c61-9e95-48567be59a3c
📒 Files selected for processing (2)
model_gateway/src/routers/grpc/regular/processor.rsmodel_gateway/src/routers/grpc/utils/message_utils.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: aef3b03628
ℹ️ 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".
…on-streaming) (#747) Signed-off-by: Simo Lin <linsimo.mark@gmail.com> Signed-off-by: VS Chandra Mourya <msrinivasa@together.ai>
Summary
Add Stage 7 (response processing) for the Messages API gRPC pipeline — converts backend
ProtoGenerateCompleteinto AnthropicMessageformat withContentBlockconstruction andStopReasonmapping. Non-streaming only; streaming deferred to follow-up PR.Fourth PR in the Messages API gRPC support series (after #739 scaffolding, #741 preparation, #744 request building).
What changed
processor.rs: Addprocess_non_streaming_messages_response()toResponseProcessor— full pipeline: token decoding → reasoning parsing → tool call parsing → content block construction (Thinking→Text→ToolUse) →StopReasonmapping (EndTurn/MaxTokens/StopSequence/ToolUse) →messages::Usagebuildingmessages/response_processing.rs: NewMessageResponseProcessingStage— extracts execution result, dispatch metadata, tokenizer, and stop decoder fromRequestContext, delegates toResponseProcessor, storesFinalResponse::Messagesmessage_utils.rs: Addget_history_tool_calls_count_messages()for counting tool use blocks in Messages API request history (for KimiK2-style tool call ID generation) + testmessages/mod.rs: Wireresponse_processingmodule withunused_importsexpectHow
Follows the same architecture as chat's response processing, adapted for Anthropic Message types:
ChatCompletionResponse+Vec<ChatChoice>Message+Vec<ContentBlock>finish_reason: StringStopReasonenumVec<ToolCall>in messageContentBlock::ToolUseblocksreasoning_content: Option<String>ContentBlock::Thinkingblockcommon::Usage(prompt/completion)messages::Usage(input/output)Key reuse:
convert_message_tool_choice()bridges Messages → Chat ToolChoice forparse_json_schema_response()parse_tool_calls()used directly for model-predicted tool parsingcollect_responses()for response collectionTest plan
cargo clippy -p smg --all-targets --all-features -- -D warningspassescargo test -p smg— all message_utils tests pass including newtest_get_history_tool_calls_count_messagescargo fmt --checkpassesRefs: #739, #741, #744
Summary by CodeRabbit
New Features
Tests