feat(gateway): wire Messages API pipeline into gRPC routers - #753
Conversation
Add pipeline factory methods and router wiring to enable /v1/messages endpoint on gRPC routers, completing the non-streaming Messages API integration. Pipeline factory (pipeline.rs): - new_messages(): dedicated Messages pipeline with MessagePreparation, MessageRequestBuilding, and MessageResponseProcessing stages, sharing worker selection, client acquisition, dispatch metadata, and request execution stages with other pipelines - new_messages_pd(): same pipeline configured for prefill-decode dual dispatch mode - execute_messages(): pipeline execution method with metrics recording and FinalResponse::Messages extraction Router wiring: - GrpcRouter: messages_pipeline field, route_messages_impl() with retry, RouterTrait::route_messages implementation - GrpcPDRouter: messages_pipeline field, route_messages_impl() with PD retry metrics (prefill + decode workers), RouterTrait::route_messages implementation Scaffolding cleanup: - Remove #[expect(dead_code)] from context.rs (for_messages, FinalResponse::Messages) now that they are used - Remove #![allow(dead_code)] from all three Messages stage files - Remove #[expect(unused_imports)] from messages/mod.rs re-exports 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 integrates the previously developed Messages API pipeline stages into the gRPC router layer, enabling functional 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
|
There was a problem hiding this comment.
Code Review
This pull request introduces comprehensive support for a new Messages API endpoint within the model_gateway service. It involves integrating CreateMessageRequest and Message types, establishing dedicated request pipelines for both regular single-worker and prefill-decode (PD) dual-dispatch modes, and implementing specific pipeline stages for message preparation, request building, and response processing. The GrpcPDRouter and GrpcRouter are updated to utilize these new message pipelines, incorporating retry logic and metrics for both operational modes. Additionally, scaffolding attributes like #[expect(dead_code)] are removed, signifying the full activation and integration of these previously placeholder components.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 519ab51bd9
ℹ️ 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 streaming = request.stream.unwrap_or(false); | ||
|
|
||
| // Record request start | ||
| Metrics::record_router_request( | ||
| metrics_labels::ROUTER_GRPC, |
There was a problem hiding this comment.
Honor
stream=true for Messages requests
execute_messages reads request.stream only for metrics, then unconditionally runs the same pipeline and serializes FinalResponse::Messages as JSON. In this commit, both new_messages and new_messages_pd are wired to MessageResponseProcessingStage (non-streaming), so /v1/messages calls with stream: true will not produce SSE events and will instead return a normal JSON payload, which breaks Anthropic streaming clients unless you explicitly reject streaming for now.
Useful? React with 👍 / 👎.
Messages API should never expose raw special tokens like <|eot_id|> in responses. Set skip_special_tokens=true in the stop sequence decoder (was false, causing EOS tokens to leak into response text). Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Caution Review failedThe pull request is closed. ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (1)
📝 WalkthroughWalkthroughThe PR implements end-to-end Messages API support in the gRPC router infrastructure by adding routing, pipeline execution, and stage wiring for CreateMessageRequest handling, while removing scaffolding dead_code attributes across messages stages. Changes
Sequence DiagramsequenceDiagram
actor Client
participant Router as gRPC Router
participant Pipeline as RequestPipeline
participant Processor as Messages<br/>Processor
participant Response
Client->>Router: CreateMessageRequest
Router->>Pipeline: execute_messages(request, headers, model_id, components)
activate Pipeline
Pipeline->>Pipeline: MessagePreparationStage
Pipeline->>Pipeline: MessageRequestBuildingStage
Pipeline->>Processor: Process request
activate Processor
Processor->>Processor: Parse reasoning (if available)
Processor-->>Pipeline: Processed response
deactivate Processor
Pipeline->>Pipeline: MessageResponseProcessingStage
Pipeline-->>Router: FinalResponse::Messages
deactivate Pipeline
Router->>Response: Apply retry/metrics
Router-->>Client: Response
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
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.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 527a338e84
ℹ️ 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 headers_cloned = headers.cloned(); | ||
| let model_id_cloned = Some(model_id.to_string()); | ||
| let components = self.shared_components.clone(); | ||
| let pipeline = &self.messages_pipeline; |
There was a problem hiding this comment.
Handle Harmony models in Messages route selection
route_messages_impl always routes through self.messages_pipeline, but route_chat_impl in this same file explicitly uses HarmonyDetector to switch Harmony/GPT-OSS models to a different pipeline. Because Harmony preparation is documented as replacing regular preparation for those models, /v1/messages requests against Harmony models are now processed with the regular Messages stages, which can yield malformed prompts or incorrect outputs for that model class; this path should either perform Harmony-aware routing or fail fast until Harmony Messages support exists.
Useful? React with 👍 / 👎.
…asoning - message_utils.rs: pass `enable_thinking: true/false` as template kwarg instead of Anthropic-style `thinking` JSON object, matching the standard HuggingFace chat template convention (e.g. Qwen3.5 `enable_thinking`) - processor.rs: always attempt reasoning parsing when a parser is available for the model, regardless of the request's `thinking` config — some models' chat templates emit thinking tokens unconditionally - pipeline.rs: reorganize messages stage imports under regular::stages block Signed-off-by: Simon Lin <simon@seekai.ai> 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/pipeline.rs`:
- Around line 750-759: execute_messages currently reads the streaming flag but
doesn't reject stream=true, causing clients to request streaming and receive
non-streaming responses; add an early guard in execute_messages (the async fn
execute_messages) that checks the streaming boolean and returns an
error::bad_request with code "streaming_not_supported" and message "Streaming is
not yet supported for the Messages API" (mirror the pattern used in
execute_chat_for_responses) so the pipeline and downstream processing (e.g.,
process_non_streaming_messages_response) are not invoked for unsupported
streaming requests.
In `@model_gateway/src/routers/grpc/regular/processor.rs`:
- Around line 490-496: The reasoning parser availability call
(utils::check_reasoning_parser_availability with self.reasoning_parser_factory
and self.configured_reasoning_parser) is fine to keep always-on for parsing, but
you must prevent emitting a ContentBlock::Thinking when the caller set
ThinkingConfig::Disabled; update the place that constructs/emits
ContentBlock::Thinking to check the request's thinking mode (e.g.,
messages_request.thinking or equivalent) and only emit the Thinking block when
the thinking config is not Disabled (leave parsing and availability logic
unchanged).
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 5d310741-7832-4f7e-a3ca-00eafadb3104
📒 Files selected for processing (3)
model_gateway/src/routers/grpc/pipeline.rsmodel_gateway/src/routers/grpc/regular/processor.rsmodel_gateway/src/routers/grpc/utils/message_utils.rs
| // Always attempt reasoning parsing when a parser is available — some models' | ||
| // chat templates emit thinking tokens regardless of the request's `thinking` config. | ||
| let reasoning_parser_available = utils::check_reasoning_parser_availability( | ||
| &self.reasoning_parser_factory, | ||
| self.configured_reasoning_parser.as_deref(), | ||
| &messages_request.model, | ||
| ); |
There was a problem hiding this comment.
Honor thinking: disabled when deciding whether to emit a thinking block.
At Line 490, reasoning parsing is now enabled regardless of request thinking mode. That is fine for cleanup, but it can cause ContentBlock::Thinking to be emitted downstream even when the caller explicitly set ThinkingConfig::Disabled.
🔧 Proposed fix (parse always, emit only when requested)
- let reasoning_parser_available = utils::check_reasoning_parser_availability(
+ let reasoning_requested = matches!(
+ messages_request.thinking.as_ref(),
+ Some(messages::ThinkingConfig::Enabled { .. })
+ );
+
+ let reasoning_parser_available = utils::check_reasoning_parser_availability(
&self.reasoning_parser_factory,
self.configured_reasoning_parser.as_deref(),
&messages_request.model,
);
@@
Ok(result) => {
- if !result.reasoning_text.is_empty() {
+ if reasoning_requested && !result.reasoning_text.is_empty() {
reasoning_text = Some(result.reasoning_text);
}
processed_text = result.normal_text;
}🤖 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 490 - 496,
The reasoning parser availability call
(utils::check_reasoning_parser_availability with self.reasoning_parser_factory
and self.configured_reasoning_parser) is fine to keep always-on for parsing, but
you must prevent emitting a ContentBlock::Thinking when the caller set
ThinkingConfig::Disabled; update the place that constructs/emits
ContentBlock::Thinking to check the request's thinking mode (e.g.,
messages_request.thinking or equivalent) and only emit the Thinking block when
the thinking config is not Disabled (leave parsing and availability logic
unchanged).
There was a problem hiding this comment.
We don't need to honor here. Our reasoning parser already handles the case inside. If a model has disabled_thinking, reasoning parser won't produce any reasoning content. This is the right way.
There was a problem hiding this comment.
@CatherineSue, understood — thanks for the clarification! Since the reasoning parser already suppresses reasoning content internally when thinking is disabled, the always-on reasoning_parser_available check is the right approach and my concern was unfounded.
✏️ Learnings added
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 753
File: model_gateway/src/routers/grpc/regular/processor.rs:490-496
Timestamp: 2026-03-13T18:45:38.016Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/grpc/regular/processor.rs: The reasoning parser (accessed via utils::check_reasoning_parser_availability / utils::get_reasoning_parser) handles ThinkingConfig::Disabled internally — it will not produce any reasoning content when the model has thinking disabled. Therefore, it is correct and intentional to always attempt reasoning parsing whenever a parser is available (reasoning_parser_available = utils::check_reasoning_parser_availability(...)), without additionally gating on the request's ThinkingConfig. Do not flag the absence of a ThinkingConfig::Enabled guard around reasoning parsing in process_non_streaming_messages_response as a bug.
Note: Learnings are effective only in the context of similar code segments. To apply general coding guidelines, please configure review instructions. You can manage existing learnings in the app.
🧠 Learnings used
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 495
File: model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs:121-131
Timestamp: 2026-02-21T18:45:58.696Z
Learning: In repo lightseekorg/smg, all multimodal processing failures in model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs currently return 400 Bad Request for simplicity, because the underlying `anyhow::Error` from the multimodal crate doesn't distinguish error types (client vs upstream failures). Error categorization (e.g., mapping upstream fetch failures to 502 Bad Gateway) is deferred to a follow-up when more failure modes need differentiation.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs:206-211
Timestamp: 2026-02-21T02:31:17.841Z
Learning: Repo: lightseekorg/smg
File: model_gateway/src/routers/grpc/client.rs
Context: GrpcClient::embed match arm
Learning: The catch-all panic in GrpcClient::embed for mismatched client/request types or unsupported embedding backends is intentional to catch invariant violations and unsupported configurations at development time. Converting this path to a returned error is out of scope for PR `#489`.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 570
File: model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs:87-97
Timestamp: 2026-03-01T05:57:56.940Z
Learning: In model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs, when multimodal content is detected, an empty tokenizer_source must be rejected with a bad_request (multimodal_config_missing) before calling process_multimodal(). Without a valid tokenizer source, get_or_load_config() fails with a confusing file-not-found error downstream. The early guard provides a clear error message. This differs from request_building.rs, where empty tokenizer_source is a valid fallback for non-multimodal config loading.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 497
File: model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs:108-113
Timestamp: 2026-02-21T23:56:04.191Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs: When fetching tokenizer_source via ctx.components.tokenizer_registry.get_by_name(model_id).map(|e| e.source).unwrap_or_default(), an empty string is an intentional valid fallback when the model isn't in the registry. This allows config loading to proceed with the default path. Comments documenting this fallback are not necessary to keep noise down.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs:189-194
Timestamp: 2026-02-21T02:32:13.396Z
Learning: Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs
Context: GrpcClient::generate fallback match arm
Learning: The catch-all panic for mismatched client/request types in GrpcClient::generate is intentional to enforce pipeline invariants; do not convert this to a returned error in PR `#489` or similar lint-only changes.
Learnt from: kzjeef
Repo: lightseekorg/smg PR: 469
File: model_gateway/src/routers/http/pd_router.rs:758-786
Timestamp: 2026-02-19T02:49:37.991Z
Learning: In model_gateway/src/routers/http/pd_router.rs, the extract_chat_request_text function intentionally concatenates chat messages without separators. This is because the radix tree performs character-level prefix matching, and separators would reduce match ratios for multi-turn conversations. No separators gives the best prefix overlap between turns for cache-aware pre-prefill routing.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 726
File: model_gateway/src/routers/openai/chat.rs:191-220
Timestamp: 2026-03-11T16:20:43.854Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/openai/chat.rs: OpenAI's Chat Completions streaming endpoint returns `text/event-stream` as the Content-Type even for error responses during streaming (mid-stream errors are encoded as SSE `error` events). Unconditionally setting `Content-Type: text/event-stream` in the `is_streaming` branch of `route_chat` is correct for OpenAI-compatible upstreams. Do not flag this as a bug for PR `#726` or similar OpenAI-router PRs.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 679
File: model_gateway/src/server.rs:247-250
Timestamp: 2026-03-10T00:16:36.169Z
Learning: In repo lightseekorg/smg, file model_gateway/src/server.rs: In the `v1_interactions` handler, `model_id` is intentionally resolved as `body.model.as_deref().or(body.agent.as_deref())`. A 400 Bad Request should only be returned when *both* `model` and `agent` are absent. Treating a missing `model` (but present `agent`) as a 400 is incorrect. Any future agent→model resolution step is separate from this extraction logic.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 468
File: model_gateway/src/routers/openai/router.rs:343-348
Timestamp: 2026-02-18T18:58:21.764Z
Learning: In `model_gateway/src/routers/openai/router.rs`, the `load_input_history` function intentionally uses graceful degradation when storage operations fail (e.g., `get_response_chain`, `list_items`). When these operations fail, a warning is logged but the request proceeds with the directly provided input rather than returning an error response. This is by design to support better UX, as requests can still succeed without history context (e.g., first message in a conversation).
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 588
File: model_gateway/src/routers/grpc/multimodal.rs:453-514
Timestamp: 2026-03-03T18:03:45.713Z
Learning: In repo lightseekorg/smg, backend assembly functions in model_gateway/src/routers/grpc/multimodal.rs (e.g., assemble_sglang, assemble_vllm, assemble_trtllm) are tested via E2E tests rather than unit tests, as unit tests for these functions are not considered worthwhile.
Learnt from: pallasathena92
Repo: lightseekorg/smg PR: 687
File: model_gateway/src/routers/openai/realtime/webrtc.rs:238-289
Timestamp: 2026-03-11T01:29:56.655Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/openai/realtime/webrtc.rs: `Metrics::record_router_request` is already emitted in router.rs (around line 1152) before `handle_realtime_webrtc` is called, and `Metrics::record_router_error` is emitted inside `handle_realtime_webrtc` for the no-workers case. The missing instrumentation is success/duration recording after `setup_and_spawn_bridge` returns — this is a metrics improvement deferred to a follow-up PR, not a correctness gap. Do not flag missing success/duration metrics as a blocking issue for PR `#687`.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 690
File: model_gateway/src/core/worker_registry.rs:667-670
Timestamp: 2026-03-10T05:04:49.809Z
Learning: In repo lightseekorg/smg, file model_gateway/src/core/worker_registry.rs: `any_external_worker_supports_model` intentionally uses `healthy_only = true` for two reasons: (1) The 503 ("service unavailable") path in `select_worker_for_model` is specifically for the circuit-breaker case — healthy workers whose circuit breaker is open — while workers failing health checks fall through to 404. (2) Unhealthy workers have stale model lists (models registered at startup may no longer be accurate) and should not be trusted for model existence checks. This design follows the upstream sglang pattern from sgl-project/sglang#15611. Do not flag `healthy_only = true` in `any_external_worker_supports_model` as a bug.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 690
File: model_gateway/src/core/worker_registry.rs:667-670
Timestamp: 2026-03-10T04:59:38.803Z
Learning: In repo lightseekorg/smg, file model_gateway/src/core/worker_registry.rs: `any_external_worker_supports_model` intentionally uses `healthy_only = true`. The 503 ("service unavailable") path in `select_worker_for_model` is specifically for the circuit-breaker case — healthy workers whose circuit breaker is open. Workers that are genuinely unhealthy (failing health checks) are intentionally excluded: the model falls through to 404 for those. This design follows the upstream sglang pattern established in sgl-project/sglang#15611 ("[model-gateway] return 503 when all workers are circuit-broken"). Do not flag the `healthy_only = true` argument in this method as a bug.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/benches/wasm_middleware_latency.rs:88-91
Timestamp: 2026-02-21T02:37:02.009Z
Learning: Repo: lightseekorg/smg — For clippy-only/enforcement PRs (e.g., PR `#489`), even micro-optimizations (like replacing an async closure with std::future::ready in benches such as model_gateway/benches/wasm_middleware_latency.rs) should be deferred to a follow-up PR rather than included inline.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 417
File: model_gateway/src/routers/gemini/router.rs:0-0
Timestamp: 2026-03-05T03:03:55.404Z
Learning: In repo lightseekorg/smg, the Gemini router (model_gateway/src/routers/gemini/router.rs) intentionally uses `mpsc::unbounded_channel` for its SSE streaming path to stay consistent with the gRPC router's streaming paths, which also use unbounded channels across the board. Backpressure handling is deferred to the dedicated streaming request implementation PR and should not be flagged as a defect in skeleton/scaffold PRs for this router.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 723
File: model_gateway/tests/api/interactions_api_test.rs:316-346
Timestamp: 2026-03-11T05:32:45.536Z
Learning: In repo lightseekorg/smg, `test_interactions_multiple_workers` in `model_gateway/tests/api/interactions_api_test.rs` intentionally only verifies that requests succeed with multiple Gemini backends (connectivity), not load distribution. Load distribution across workers is covered separately in `tests/routing/load_balancing_test.rs`. Do not flag the absence of per-worker request count assertions in this test as a gap — distribution verification is deferred to a follow-up PR once the Gemini router implementation matures.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 447
File: model_gateway/src/routers/grpc/client.rs:312-328
Timestamp: 2026-02-17T20:30:27.647Z
Learning: Actionable guideline: In model_gateway gRPC metadata discovery (specifically in model_gateway/src/routers/grpc/...), verify how keys are handled for different proto sources. SGLang uses short-form keys (tp_size, dp_size, pp_size) via pick_prost_fields() without normalization, while vLLM/TRT-LLM use long-form keys (tensor_parallel_size, pipeline_parallel_size) that pass through flat_labels() and are normalized by normalize_grpc_keys() in discover_metadata.rs after model_info.to_labels() and device/server_info.to_labels(). Ensure reviewers check that the code paths correctly reflect these normalization rules and that tests cover both code paths.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 399
File: protocols/src/interactions.rs:505-509
Timestamp: 2026-02-19T03:08:50.192Z
Learning: In code reviews for Rust projects using the validator crate (v0.20.0), ensure that custom validation functions for numeric primitive types (e.g., f32, i32, u32, i16, etc.) accept the value by value, not by reference. Example: fn validate(value: f32) { ... }. The validator derive macro has a hardcoded list of numeric types that are passed by value, while all other types are passed by reference. Apply this guideline whenever validating numeric fields to align with the derive macro behavior.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/core/token_bucket.rs:58-63
Timestamp: 2026-02-21T02:30:51.443Z
Learning: For lint-only/Clippy enforcement PRs in this repository, avoid introducing behavioral changes (e.g., new input validation or logic changes). Treat such PRs as non-functional changes and plan a separate follow-up issue/PR for hardening or behavior changes. This applies broadly to Rust files across the repo; during review, focus on lint/style corrections and clearly note any intentional exceptions.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/responses.rs:928-931
Timestamp: 2026-02-21T02:36:00.882Z
Learning: In Rust code across the repository, use the marker INVARIANT: to document assumptions in safe code. Reserve SAFETY: for explaining why unsafe blocks are sound. This improves clarity of invariants and safety reasoning. Example reference: protocols/src/responses.rs near validate_tool_choice_with_tools().
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: mesh/src/sync.rs:83-83
Timestamp: 2026-02-21T02:37:01.416Z
Learning: General Rust formatting rule: format! with implicit captures only supports simple identifiers, not full expressions like {state.model_id}. For cases where you want to interpolate a field or expression, bind the value first and interpolate the binding, e.g., let model_id = &state.model_id; and then use format!("policy:{}", model_id). In the specific file mesh/src/sync.rs, prefer format!("policy:{}", state.model_id) or bind to a local variable if you need named interpolation, to keep clarity and avoid unintended captures.
There was a problem hiding this comment.
For instance, for a model that doesn't support enable_thinking, meaning it always has thinking content. Then the model itself won't respect users setting. Having your guard here would make reasoning content always inside regular text block. It would be wrong. Handling this in reasoning parser is correct.
There was a problem hiding this comment.
@CatherineSue, that's a great clarification — thank you! The concrete case makes the design intent very clear: for models that unconditionally emit thinking tokens (because they don't honour enable_thinking), placing a ThinkingConfig::Enabled guard at the emission site would silently swallow all that thinking content into the regular text block. Delegating the decision entirely to the reasoning parser, which has model-level awareness, is the correct approach.
I'll update the learning with this additional context.
✏️ Learnings added
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 753
File: model_gateway/src/routers/grpc/regular/processor.rs:490-496
Timestamp: 2026-03-13T18:47:23.096Z
Learning: In model_gateway/src/routers/grpc/regular/processor.rs, keep attempting reasoning parsing whenever a reasoning_parser is available (i.e., when reasoning_parser_available = utils::check_reasoning_parser_availability(...)) regardless of the request's ThinkingConfig. This is correct for two reasons: (1) The reasoning parser handles ThinkingConfig::Disabled internally and will not produce reasoning content when thinking is disabled. (2) Some models do not support enable_thinking and will always emit thinking tokens regardless of the user's ThinkingConfig setting; placing a ThinkingConfig::Enabled guard at the ContentBlock::Thinking emission site would bury that thinking content inside the regular text block instead of properly separating it. The reasoning parser, having model-level awareness, is the right place to handle this. Do not flag the absence of a ThinkingConfig::Enabled guard around reasoning parsing or ContentBlock::Thinking emission in process_non_streaming_messages_response as a bug. This decision is file-specific and should not be generalized beyond this processor.rs implementation.
Note: Learnings are effective only in the context of similar code segments. To apply general coding guidelines, please configure review instructions. You can manage existing learnings in the app.
🧠 Learnings used
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 753
File: model_gateway/src/routers/grpc/regular/processor.rs:490-496
Timestamp: 2026-03-13T18:45:38.016Z
Learning: In model_gateway/src/routers/grpc/regular/processor.rs, keep attempting reasoning parsing whenever a reasoning_parser is available (i.e., when reasoning_parser_available = utils::check_reasoning_parser_availability(...)) regardless of the request's ThinkingConfig. Reasoning content is suppressed when ThinkingConfig is Disabled, so gating on ThinkingConfig.Enabled is unnecessary. Do not flag the absence of a ThinkingConfig::Enabled guard around reasoning parsing in process_non_streaming_messages_response as a bug. This decision is file-specific and should not be generalized beyond this processor.rs implementation.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 495
File: model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs:121-131
Timestamp: 2026-02-21T18:45:58.696Z
Learning: In repo lightseekorg/smg, all multimodal processing failures in model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs currently return 400 Bad Request for simplicity, because the underlying `anyhow::Error` from the multimodal crate doesn't distinguish error types (client vs upstream failures). Error categorization (e.g., mapping upstream fetch failures to 502 Bad Gateway) is deferred to a follow-up when more failure modes need differentiation.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs:206-211
Timestamp: 2026-02-21T02:31:17.841Z
Learning: Repo: lightseekorg/smg
File: model_gateway/src/routers/grpc/client.rs
Context: GrpcClient::embed match arm
Learning: The catch-all panic in GrpcClient::embed for mismatched client/request types or unsupported embedding backends is intentional to catch invariant violations and unsupported configurations at development time. Converting this path to a returned error is out of scope for PR `#489`.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 570
File: model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs:87-97
Timestamp: 2026-03-01T05:57:56.940Z
Learning: In model_gateway/src/routers/grpc/regular/stages/chat/preparation.rs, when multimodal content is detected, an empty tokenizer_source must be rejected with a bad_request (multimodal_config_missing) before calling process_multimodal(). Without a valid tokenizer source, get_or_load_config() fails with a confusing file-not-found error downstream. The early guard provides a clear error message. This differs from request_building.rs, where empty tokenizer_source is a valid fallback for non-multimodal config loading.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 497
File: model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs:108-113
Timestamp: 2026-02-21T23:56:04.191Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs: When fetching tokenizer_source via ctx.components.tokenizer_registry.get_by_name(model_id).map(|e| e.source).unwrap_or_default(), an empty string is an intentional valid fallback when the model isn't in the registry. This allows config loading to proceed with the default path. Comments documenting this fallback are not necessary to keep noise down.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs:189-194
Timestamp: 2026-02-21T02:32:13.396Z
Learning: Repo: lightseekorg/smg PR: 489
File: model_gateway/src/routers/grpc/client.rs
Context: GrpcClient::generate fallback match arm
Learning: The catch-all panic for mismatched client/request types in GrpcClient::generate is intentional to enforce pipeline invariants; do not convert this to a returned error in PR `#489` or similar lint-only changes.
Learnt from: kzjeef
Repo: lightseekorg/smg PR: 469
File: model_gateway/src/routers/http/pd_router.rs:758-786
Timestamp: 2026-02-19T02:49:37.991Z
Learning: In model_gateway/src/routers/http/pd_router.rs, the extract_chat_request_text function intentionally concatenates chat messages without separators. This is because the radix tree performs character-level prefix matching, and separators would reduce match ratios for multi-turn conversations. No separators gives the best prefix overlap between turns for cache-aware pre-prefill routing.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 726
File: model_gateway/src/routers/openai/chat.rs:191-220
Timestamp: 2026-03-11T16:20:43.854Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/openai/chat.rs: OpenAI's Chat Completions streaming endpoint returns `text/event-stream` as the Content-Type even for error responses during streaming (mid-stream errors are encoded as SSE `error` events). Unconditionally setting `Content-Type: text/event-stream` in the `is_streaming` branch of `route_chat` is correct for OpenAI-compatible upstreams. Do not flag this as a bug for PR `#726` or similar OpenAI-router PRs.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 679
File: model_gateway/src/server.rs:247-250
Timestamp: 2026-03-10T00:16:36.169Z
Learning: In repo lightseekorg/smg, file model_gateway/src/server.rs: In the `v1_interactions` handler, `model_id` is intentionally resolved as `body.model.as_deref().or(body.agent.as_deref())`. A 400 Bad Request should only be returned when *both* `model` and `agent` are absent. Treating a missing `model` (but present `agent`) as a 400 is incorrect. Any future agent→model resolution step is separate from this extraction logic.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 468
File: model_gateway/src/routers/openai/router.rs:343-348
Timestamp: 2026-02-18T18:58:21.764Z
Learning: In `model_gateway/src/routers/openai/router.rs`, the `load_input_history` function intentionally uses graceful degradation when storage operations fail (e.g., `get_response_chain`, `list_items`). When these operations fail, a warning is logged but the request proceeds with the directly provided input rather than returning an error response. This is by design to support better UX, as requests can still succeed without history context (e.g., first message in a conversation).
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 690
File: model_gateway/src/core/worker_registry.rs:667-670
Timestamp: 2026-03-10T05:04:49.809Z
Learning: In repo lightseekorg/smg, file model_gateway/src/core/worker_registry.rs: `any_external_worker_supports_model` intentionally uses `healthy_only = true` for two reasons: (1) The 503 ("service unavailable") path in `select_worker_for_model` is specifically for the circuit-breaker case — healthy workers whose circuit breaker is open — while workers failing health checks fall through to 404. (2) Unhealthy workers have stale model lists (models registered at startup may no longer be accurate) and should not be trusted for model existence checks. This design follows the upstream sglang pattern from sgl-project/sglang#15611. Do not flag `healthy_only = true` in `any_external_worker_supports_model` as a bug.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 690
File: model_gateway/src/core/worker_registry.rs:667-670
Timestamp: 2026-03-10T04:59:38.803Z
Learning: In repo lightseekorg/smg, file model_gateway/src/core/worker_registry.rs: `any_external_worker_supports_model` intentionally uses `healthy_only = true`. The 503 ("service unavailable") path in `select_worker_for_model` is specifically for the circuit-breaker case — healthy workers whose circuit breaker is open. Workers that are genuinely unhealthy (failing health checks) are intentionally excluded: the model falls through to 404 for those. This design follows the upstream sglang pattern established in sgl-project/sglang#15611 ("[model-gateway] return 503 when all workers are circuit-broken"). Do not flag the `healthy_only = true` argument in this method as a bug.
Learnt from: pallasathena92
Repo: lightseekorg/smg PR: 687
File: model_gateway/src/routers/openai/realtime/webrtc.rs:238-289
Timestamp: 2026-03-11T01:29:56.655Z
Learning: In repo lightseekorg/smg, file model_gateway/src/routers/openai/realtime/webrtc.rs: `Metrics::record_router_request` is already emitted in router.rs (around line 1152) before `handle_realtime_webrtc` is called, and `Metrics::record_router_error` is emitted inside `handle_realtime_webrtc` for the no-workers case. The missing instrumentation is success/duration recording after `setup_and_spawn_bridge` returns — this is a metrics improvement deferred to a follow-up PR, not a correctness gap. Do not flag missing success/duration metrics as a blocking issue for PR `#687`.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 495
File: model_gateway/src/routers/grpc/multimodal.rs:153-166
Timestamp: 2026-02-21T18:39:37.571Z
Learning: In repo lightseekorg/smg, the gRPC multimodal pipeline in model_gateway/src/routers/grpc/multimodal.rs intentionally skips ContentPart::VideoUrl silently (without warning) because video support is not yet implemented; adding warnings would be noisy for legitimate requests containing video parts that aren't processed yet. Proper video handling will be added when the feature is implemented.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/benches/wasm_middleware_latency.rs:88-91
Timestamp: 2026-02-21T02:37:02.009Z
Learning: Repo: lightseekorg/smg — For clippy-only/enforcement PRs (e.g., PR `#489`), even micro-optimizations (like replacing an async closure with std::future::ready in benches such as model_gateway/benches/wasm_middleware_latency.rs) should be deferred to a follow-up PR rather than included inline.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 417
File: model_gateway/src/routers/gemini/router.rs:0-0
Timestamp: 2026-03-05T03:03:55.404Z
Learning: In repo lightseekorg/smg, the Gemini router (model_gateway/src/routers/gemini/router.rs) intentionally uses `mpsc::unbounded_channel` for its SSE streaming path to stay consistent with the gRPC router's streaming paths, which also use unbounded channels across the board. Backpressure handling is deferred to the dedicated streaming request implementation PR and should not be flagged as a defect in skeleton/scaffold PRs for this router.
Learnt from: CatherineSue
Repo: lightseekorg/smg PR: 588
File: model_gateway/src/routers/grpc/multimodal.rs:453-514
Timestamp: 2026-03-03T18:03:45.713Z
Learning: In repo lightseekorg/smg, backend assembly functions in model_gateway/src/routers/grpc/multimodal.rs (e.g., assemble_sglang, assemble_vllm, assemble_trtllm) are tested via E2E tests rather than unit tests, as unit tests for these functions are not considered worthwhile.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 723
File: model_gateway/tests/api/interactions_api_test.rs:316-346
Timestamp: 2026-03-11T05:32:45.536Z
Learning: In repo lightseekorg/smg, `test_interactions_multiple_workers` in `model_gateway/tests/api/interactions_api_test.rs` intentionally only verifies that requests succeed with multiple Gemini backends (connectivity), not load distribution. Load distribution across workers is covered separately in `tests/routing/load_balancing_test.rs`. Do not flag the absence of per-worker request count assertions in this test as a gap — distribution verification is deferred to a follow-up PR once the Gemini router implementation matures.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 447
File: model_gateway/src/routers/grpc/client.rs:312-328
Timestamp: 2026-02-17T20:30:27.647Z
Learning: Actionable guideline: In model_gateway gRPC metadata discovery (specifically in model_gateway/src/routers/grpc/...), verify how keys are handled for different proto sources. SGLang uses short-form keys (tp_size, dp_size, pp_size) via pick_prost_fields() without normalization, while vLLM/TRT-LLM use long-form keys (tensor_parallel_size, pipeline_parallel_size) that pass through flat_labels() and are normalized by normalize_grpc_keys() in discover_metadata.rs after model_info.to_labels() and device/server_info.to_labels(). Ensure reviewers check that the code paths correctly reflect these normalization rules and that tests cover both code paths.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 399
File: protocols/src/interactions.rs:505-509
Timestamp: 2026-02-19T03:08:50.192Z
Learning: In code reviews for Rust projects using the validator crate (v0.20.0), ensure that custom validation functions for numeric primitive types (e.g., f32, i32, u32, i16, etc.) accept the value by value, not by reference. Example: fn validate(value: f32) { ... }. The validator derive macro has a hardcoded list of numeric types that are passed by value, while all other types are passed by reference. Apply this guideline whenever validating numeric fields to align with the derive macro behavior.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/core/token_bucket.rs:58-63
Timestamp: 2026-02-21T02:30:51.443Z
Learning: For lint-only/Clippy enforcement PRs in this repository, avoid introducing behavioral changes (e.g., new input validation or logic changes). Treat such PRs as non-functional changes and plan a separate follow-up issue/PR for hardening or behavior changes. This applies broadly to Rust files across the repo; during review, focus on lint/style corrections and clearly note any intentional exceptions.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/responses.rs:928-931
Timestamp: 2026-02-21T02:36:00.882Z
Learning: In Rust code across the repository, use the marker INVARIANT: to document assumptions in safe code. Reserve SAFETY: for explaining why unsafe blocks are sound. This improves clarity of invariants and safety reasoning. Example reference: protocols/src/responses.rs near validate_tool_choice_with_tools().
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: mesh/src/sync.rs:83-83
Timestamp: 2026-02-21T02:37:01.416Z
Learning: General Rust formatting rule: format! with implicit captures only supports simple identifiers, not full expressions like {state.model_id}. For cases where you want to interpolate a field or expression, bind the value first and interpolate the binding, e.g., let model_id = &state.model_id; and then use format!("policy:{}", model_id). In the specific file mesh/src/sync.rs, prefer format!("policy:{}", state.model_id) or bind to a local variable if you need named interpolation, to keep clarity and avoid unintended captures.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 26b7cedfb2
ℹ️ 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: 1
🤖 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/reasoning_parser/src/parsers/base.rs`:
- Around line 55-64: The debug log is slicing the string by bytes with
&text[..text.len().min(100)] which can panic on UTF-8 boundaries; change the log
to produce a character-safe preview by truncating to 100 characters instead
(e.g., build a short preview via text.chars().take(100).collect::<String>() or
use text.floor_char_boundary(100) on Rust 1.80+) and pass that preview into the
tracing::debug call (references: self.model_type, in_reasoning,
text_contains_start, text_contains_end, text_len, text_preview).
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 32f1efcf-34da-41c6-91be-5db0876413d9
📒 Files selected for processing (2)
crates/reasoning_parser/Cargo.tomlcrates/reasoning_parser/src/parsers/base.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 7ef3c89d48
ℹ️ 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.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 93de5a9da0
ℹ️ 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".
Different model chat templates use different kwarg names for controlling thinking/reasoning mode. Qwen3 uses `enable_thinking` while Kimi-K2.5 uses `thinking`. Pass both so the correct one is picked up regardless of which template is loaded. - model_gateway/src/routers/grpc/utils/message_utils.rs: insert both `enable_thinking` and `thinking` booleans into template kwargs Refs: #753 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Add streaming support (stream: true) for the Messages API in the gRPC
router, emitting Anthropic SSE events (event: {type}\ndata: {json}\n\n).
What changed:
- streaming.rs: Add process_messages_streaming_response entry point,
process_messages_streaming_chunks core loop, and
process_dual_messages_streaming_chunks for PD mode. Add SSE helpers
(format_messages_sse_into, send_messages_event, message_event_type_name)
using reusable buffer pattern. Add process_messages_reasoning helper.
Implement incremental tool call streaming via parse_incremental with
content block state machine (thinking/text/tool_use blocks).
- response_processing.rs: Add StreamingProcessor to
MessageResponseProcessingStage. Streaming branch creates SSE response
and attaches load guards; non-streaming branch unchanged.
- pipeline.rs: Wire StreamingProcessor into both new_messages and
new_messages_pd pipelines by cloning parser factories.
Design decisions:
- Architecture matches chat streaming exactly: same StreamingProcessor
struct, same entry point pattern, same reusable SSE buffer, same
incremental tool parser via parse_incremental/get_unstreamed_tool_args.
- No HashMap per-index since Messages API is always n=1.
- Specific function (ToolChoice::Tool) streams arguments directly;
regular/required modes use incremental parser.
- Error events use MessageStreamEvent::Error (no data: [DONE]).
Refs: #753
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
When tool_choice is Any or Tool, the model's entire output is tool call JSON. The reasoning parser was misclassifying this JSON as thinking content, leaving nothing for the tool parser — resulting in no tool_use blocks and stop_reason: end_turn instead of tool_use. What changed: - processor.rs: Move used_json_schema computation before Step 1 (reasoning parsing). Skip reasoning parser when used_json_schema is true, since the model output is constrained to tool call JSON. Remove duplicate used_json_schema definition from Step 2. Refs: #753 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Add streaming support (stream: true) for the Messages API in the gRPC
router, emitting Anthropic SSE events (event: {type}\ndata: {json}\n\n).
What changed:
- streaming.rs: Add process_messages_streaming_response entry point,
process_messages_streaming_chunks core loop, and
process_dual_messages_streaming_chunks for PD mode. Add SSE helpers
(format_messages_sse_into, send_messages_event, message_event_type_name)
using reusable buffer pattern. Add process_messages_reasoning helper.
Implement incremental tool call streaming via parse_incremental with
content block state machine (thinking/text/tool_use blocks).
- response_processing.rs: Add StreamingProcessor to
MessageResponseProcessingStage. Streaming branch creates SSE response
and attaches load guards; non-streaming branch unchanged.
- pipeline.rs: Wire StreamingProcessor into both new_messages and
new_messages_pd pipelines by cloning parser factories.
Design decisions:
- Architecture matches chat streaming exactly: same StreamingProcessor
struct, same entry point pattern, same reusable SSE buffer, same
incremental tool parser via parse_incremental/get_unstreamed_tool_args.
- No HashMap per-index since Messages API is always n=1.
- Specific function (ToolChoice::Tool) streams arguments directly;
regular/required modes use incremental parser.
- Error events use MessageStreamEvent::Error (no data: [DONE]).
Refs: #753
Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
When tool_choice is Any or Tool, the model's entire output is tool call JSON. The reasoning parser was misclassifying this JSON as thinking content, leaving nothing for the tool parser — resulting in no tool_use blocks and stop_reason: end_turn instead of tool_use. What changed: - processor.rs: Move used_json_schema computation before Step 1 (reasoning parsing). Skip reasoning parser when used_json_schema is true, since the model output is constrained to tool call JSON. Remove duplicate used_json_schema definition from Step 2. Refs: #753 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
…format Add messages_streaming_test.rs covering Anthropic Messages API spec conformance: - MessageStreamEvent serialization for all 8 event variants - Non-streaming Message response golden tests (text, tool_use, thinking) - StopReason serialization/deserialization round-trip - SSE wire format validation (event: + data: + double newline) - Complete streaming event sequence lifecycle (text + tool_use) - stream field default handling (true/false/omitted) - message_delta stop_reason variants (stop_sequence, max_tokens) All test data references: https://docs.anthropic.com/en/api/messages Refs: #753
…son resolver
Extract build_messages_content_blocks() and resolve_messages_stop_reason()
from process_non_streaming_messages_response() into testable functions.
Tests cover:
- Content block ordering: thinking → text → tool_use
- Empty text omission
- Invalid tool arguments fallback to {}
- Multiple tool calls
- stop_reason priority: tool_calls > stop_sequence > length > end_turn
- tool_calls flag takes priority over matched stop_sequence
Refs: #753
…format Add messages_streaming_test.rs covering Anthropic Messages API spec conformance: - MessageStreamEvent serialization for all 8 event variants - Non-streaming Message response golden tests (text, tool_use, thinking) - StopReason serialization/deserialization round-trip - SSE wire format validation (event: + data: + double newline) - Complete streaming event sequence lifecycle (text + tool_use) - stream field default handling (true/false/omitted) - message_delta stop_reason variants (stop_sequence, max_tokens) All test data references: https://docs.anthropic.com/en/api/messages Refs: #753
…format Add messages_streaming_test.rs covering Anthropic Messages API spec conformance: - MessageStreamEvent serialization for all 8 event variants - Non-streaming Message response golden tests (text, tool_use, thinking) - StopReason serialization/deserialization round-trip - SSE wire format validation (event: + data: + double newline) - Complete streaming event sequence lifecycle (text + tool_use) - stream field default handling (true/false/omitted) - message_delta stop_reason variants (stop_sequence, max_tokens) All test data references: https://docs.anthropic.com/en/api/messages Refs: #753
Add e2e_test/messages/test_grpc_messages.py with 8 tests covering: Non-streaming (TestGrpcMessagesBasic): - Basic text response (structure + fields) - System prompt - Multi-turn conversation - stream=false returns JSON not SSE Streaming (TestGrpcMessagesStreaming): - SSE event sequence completeness - Text delta concatenation Tool use (TestGrpcMessagesToolUse): - Non-streaming tool_use with tool_choice=tool - Streaming input_json_delta events Uses setup_backend=["grpc"] fixture for CI compatibility. Model: Llama-3.1-8B-Instruct (same as chat_completions tests). Refs: #753, #758
Add e2e_test/messages/test_grpc_messages.py with 8 tests covering: Non-streaming (TestGrpcMessagesBasic): - Basic text response (structure + fields) - System prompt - Multi-turn conversation - stream=false returns JSON not SSE Streaming (TestGrpcMessagesStreaming): - SSE event sequence completeness - Text delta concatenation Tool use (TestGrpcMessagesToolUse): - Non-streaming tool_use with tool_choice=tool - Streaming input_json_delta events Uses setup_backend=["grpc"] fixture for CI compatibility. Model: Llama-3.1-8B-Instruct (same as chat_completions tests). Refs: #753, #758
Add e2e_test/messages/test_grpc_messages.py with 8 tests covering: Non-streaming (TestGrpcMessagesBasic): - Basic text response (structure + fields) - System prompt - Multi-turn conversation - stream=false returns JSON not SSE Streaming (TestGrpcMessagesStreaming): - SSE event sequence completeness - Text delta concatenation Tool use (TestGrpcMessagesToolUse): - Non-streaming tool_use with tool_choice=tool - Streaming input_json_delta events Uses setup_backend=["grpc"] fixture for CI compatibility. Model: Llama-3.1-8B-Instruct (same as chat_completions tests). Refs: #753, #758
Add e2e_test/messages to the e2e-1gpu-chat job's test_dirs so that Messages API gRPC router tests run alongside chat_completions tests. Marker-based filtering ensures: - engine(sglang,vllm) + gpu(1) tests run in the GPU job - vendor(anthropic) + gpu(0) tests remain in the CPU vendor job - No cross-contamination between the two groups Refs: #753, #758
Signed-off-by: Simo Lin <linsimo.mark@gmail.com> Signed-off-by: VS Chandra Mourya <msrinivasa@together.ai>
Summary
Wire the Messages API pipeline into gRPC routers, enabling
/v1/messagesendpoint support for both regular and PD (prefill-decode) modes. This is PR 5 in the Messages API gRPC series.Closes the pipeline factory + router wiring step from the design doc.
What changed
Pipeline factory (
pipeline.rs):new_messages()— dedicated Messages pipeline (single-worker) withMessagePreparationStage, shared middleware stages (worker selection, client acquisition, dispatch metadata, request execution), andMessageResponseProcessingStagenew_messages_pd()— same pipeline configured for PD dual dispatch (WorkerSelectionMode::PrefillDecode,ExecutionMode::DualDispatch, PD metadata injection)execute_messages()— pipeline execution with metrics recording (ENDPOINT_MESSAGES) andFinalResponse::MessagesextractionCreateMessageRequestand Messages stage typesRouter wiring (
router.rs):messages_pipelinefield onGrpcRouter, created innew()vianew_messages()route_messages_impl()withRetryExecutorretry wrapperRouterTrait::route_messagesoverride (was returning 501 Not Implemented)PD Router wiring (
pd_router.rs):messages_pipelinefield onGrpcPDRouter, created innew()vianew_messages_pd()route_messages_impl()with PD retry metrics (prefill + decode workers)RouterTrait::route_messagesoverrideScaffolding cleanup:
#[expect(dead_code)]fromcontext.rs(for_messages,FinalResponse::Messages)#![allow(dead_code)]from all three Messages stage files (preparation.rs,request_building.rs,response_processing.rs)#[expect(unused_imports)]frommessages/mod.rsre-exportsmessages/mod.rsdoc commentWhy
Previous PRs (#739, #741, #744, #747) built the individual Messages pipeline stages but left them unwired. This PR composes them into a functional pipeline and connects it to the router layer, making
/v1/messagesactually reachable on gRPC routers.How
Follows the same dedicated-pipeline pattern as embeddings/classify — Messages gets its own pipeline instance with endpoint-specific stages at positions 1, 4, 7 and shared stages at positions 2, 3, 5, 6. This avoids modifying the delegating stages (which handle Chat + Generate) and keeps the architecture clean.
Both
GrpcRouter(regular) andGrpcPDRouter(prefill-decode) get their ownmessages_pipelinefield androute_messages_impl()method, following the exact same retry + metrics patterns asroute_chat_impl().Test plan
cargo clippy -p smg --all-targets --all-features -- -D warnings— cleancargo clippy -p smg-grpc-client --all-targets --all-features -- -D warnings— cleancargo fmt --check— cleancargo test -p smg -- grpc— passing/v1/messagesrequest to gRPC-backed model (see PR description for sample request)Sample request
Prior PRs in series
RequestType::Messages,FinalResponse::Messages)MessagePreparationStage+ message_utilsMessageRequestBuildingStage+ backend sampling paramsMessageResponseProcessingStage(non-streaming)Remaining work
MessageStreamingProcessor, SSE events)Summary by CodeRabbit
Release Notes
New Features
Improvements
Chores