Skip to content

feat(completions): add Completion API streaming support to gRPC router - #978

Merged
slin1237 merged 11 commits into
mainfrom
mourya/cmp-5
Apr 1, 2026
Merged

slin1237 merged 11 commits into
mainfrom
mourya/cmp-5

Conversation

@vschandramourya

@vschandramourya vschandramourya commented Mar 30, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Add streaming support (stream: true) for the Completion API in the gRPC router,
emitting OpenAI-compatible SSE events (data: {json}\n\ndata: [DONE]\n\n).
Supports both Regular and PD (prefill-decode) execution modes.

What changed

streaming.rs (+340 lines):

  • process_completion_streaming_response: entry point — creates channel, spawns
    background task, returns SSE response. Handles Single, Dual, and Embedding (error)
    execution modes
  • process_completion_streaming_chunks: core streaming loop with per-index stop
    decoder tracking, echo handling (prepend prompt in first chunk), suffix
    handling (append after final chunk), and finish reason parsing
  • process_dual_completion_streaming_chunks: PD mode — consumes prefill stream,
    delegates decode stream to process_completion_streaming_chunks
  • format_completion_sse_into: SSE helper using reusable buffer pattern

response_processing.rs (rewritten):

  • CompletionResponseProcessingStage now holds Arc<StreamingProcessor>
  • execute() checks ctx.is_streaming() first: streaming branch calls
    process_completion_streaming_response and attaches load guards via
    AttachedBody; non-streaming branch unchanged

pipeline.rs (+26 lines):

  • Both new_completion() and new_completion_pd() now create a
    StreamingProcessor and pass it into CompletionResponseProcessingStage::new()

How

Follows the same streaming architecture used by other gRPC endpoints in this repo:
same StreamingProcessor struct, same entry point / chunks / dual pattern, same
reusable SSE buffer, same data: {json}\n\ndata: [DONE]\n\n wire format.

Test plan

  • cargo clippy -p smg --all-targets --all-features -- -D warnings — clean
  • cargo fmt --check — clean
  • Manual E2E: verified on GLM-5 (4×GPU, SGLang backend):
    • Non-streaming: basic, n=2, echo+suffix, stop sequences ✓
    • Streaming: basic token-by-token SSE ✓
    • Streaming: echo (prompt prepended in first chunk) ✓
    • Streaming: suffix (appended after final chunk) ✓
    • Streaming: stop sequences (finish_reason: stop + data: [DONE]) ✓

Prior PRs in series

  1. feat(completions): add native gRPC pipeline typing for /v1/completions #840 — Scaffolding
  2. feat(completions): add CompletionPreparationStage for gRPC pipeline #907 — Stage 1: CompletionPreparationStage
  3. feat(completions): add CompletionRequestBuildingStage and backend sampling params #915 — Stage 4: CompletionRequestBuildingStage
  4. feat(completions): add CompletionResponseProcessingStage for non-streaming responses #953 — Stage 7: CompletionResponseProcessingStage (non-streaming)
  5. feat(completions): wire Completion API pipeline into gRPC routers #964 — Pipeline factory + Router wiring
  6. This PR — Streaming support

Summary by CodeRabbit

  • New Features

    • SSE streaming for completions: chunked token events, guaranteed terminal [DONE] frame, and dual/prefill streaming support.
  • Improvements

    • Streaming responses delivered immediately while non-streaming flows remain unchanged.
    • Per-choice stop/echo/suffix handling, explicit finish reasons, first-token latency and token metrics.
    • Backend errors forwarded as SSE error frames; embeddings streaming returns immediate invalid-request + [DONE].

@github-actions github-actions Bot added grpc gRPC client and router changes model-gateway Model gateway crate changes labels Mar 30, 2026
@coderabbitai

coderabbitai Bot commented Mar 30, 2026 •

Copy link
Copy Markdown

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review
📝 Walkthrough

Walkthrough

Wires a shared Arc<streaming::StreamingProcessor> into completion pipelines and updates the completion response stage to route streaming requests to new SSE streaming handlers while preserving existing non‑streaming behavior.

Changes

Cohort / File(s) Summary
Pipeline Construction
model_gateway/src/routers/grpc/pipeline.rs
Allocate Arc<streaming::StreamingProcessor> (with ToolParserFactory::default(), ReasoningParserFactory::default(), None, and backend label) in RequestPipeline::new_completion and new_completion_pd, and pass it into CompletionResponseProcessingStage::new(...).
Response Processing Stage
model_gateway/src/routers/grpc/regular/stages/completion/response_processing.rs
Add streaming_processor: Arc<streaming::StreamingProcessor> field and constructor arg; execute now branches on ctx.is_streaming() to invoke streaming_processor.process_completion_streaming_response(...) for streaming (early-returning an SSE Response) or continue existing non‑streaming flow that sets ctx.state.response.final_response.
Streaming Implementation
model_gateway/src/routers/grpc/regular/streaming.rs
New SSE /v1/completions streaming support: process_completion_streaming_response, process_completion_streaming_chunks, process_dual_completion_streaming_chunks, and format_completion_sse_into. Implements token decoding, per-choice stop decoding, echo/suffix handling, prefill draining for dual streams, TTFT/metrics, backend error-as-SSE handling, and final [DONE] emission.

Sequence Diagram

sequenceDiagram
    participant Client
    participant ResponseStage as CompletionResponseProcessingStage
    participant StreamingProc as StreamingProcessor
    participant SSE as "SSE Channel"
    participant Decoder as "Token Decoder"
    participant Tokenizer

    Client->>ResponseStage: Completion request (is_streaming = true)
    ResponseStage->>StreamingProc: process_completion_streaming_response(execution_result, request, dispatch, tokenizer)
    StreamingProc->>SSE: create SSE channel (tx/rx) and return Response
    StreamingProc->>Decoder: spawn decoding task (single or dual stream)
    SSE-->>Client: SSE connection established
    Decoder->>Tokenizer: decode token IDs -> text chunk
    Tokenizer-->>Decoder: decoded text
    Decoder->>SSE: emit CompletionStreamResponse chunk
    loop until backend Complete or Error
        Decoder->>Tokenizer: next tokens
        Tokenizer-->>Decoder: chunk
        Decoder->>SSE: emit chunk
    end
    Decoder->>SSE: emit data: [DONE]
    SSE-->>Client: data: [DONE]
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly related PRs

Suggested reviewers

  • CatherineSue
  • key4ng

Poem

🐰 I nibble bytes and hum with glee,
Streams of tokens hop to me,
Chunks hop out in bright array,
SSE hums them on their way,
Hooray — completions stream with tea! 🍵

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title 'feat(completions): add Completion API streaming support to gRPC router' accurately and specifically describes the main change: implementing streaming support for the Completion API in the gRPC router.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch mourya/cmp-5

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands and usage tips.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 8ee624e063

ℹ️ 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".

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review

This pull request implements streaming support for the /v1/completions API, integrating a StreamingProcessor into the gRPC pipeline to handle SSE responses for both regular and dual-dispatch execution modes. The feedback identifies critical bugs in the multi-choice (n > 1) logic, where a single stopped index could prematurely terminate the entire stream or result in missing finish reasons for other choices. Additionally, improvements were suggested to optimize performance by avoiding unnecessary string clones and re-allocations within the streaming loop.

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 `@model_gateway/src/routers/grpc/regular/streaming.rs`:
- Around line 2326-2331: The current branch sets *is_first = false in both the
if and else branches redundantly; simplify by removing the else branch and
ensuring *is_first is set to false exactly once after handling the echo case:
when *is_first && echo, prepend prompt_text to chunk_text (using prompt_text and
chunk_text) then clear *is_first, otherwise just clear *is_first without a
separate else branch. Update the code around the is_first/echo check in the
streaming logic (the block manipulating is_first, echo, chunk_text, and
prompt_text) accordingly.
🪄 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: c7c64800-8b9d-4363-a510-db8772544952

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd4a0a and 8ee624e.

📒 Files selected for processing (3)
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/response_processing.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: e0da350cd8

ℹ️ 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".

Comment on lines +2336 to +2339
if *is_first {
if echo {
chunk_text = format!("{prompt_text}{chunk_text}");
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Handle echo when a stream emits Complete before any Chunk

Echo insertion currently happens only in the Chunk path (if *is_first { ... }), so if a backend emits Complete without sending any chunk first (for example max_tokens=0 or immediate stop), stream: true with echo: true produces no echoed prompt text at all. This diverges from non-streaming completion behavior, which always prepends the prompt when echo is enabled.

Useful? React with 👍 / 👎.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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/streaming.rs`:
- Around line 2336-2341: The code currently only prepends prompt_text when the
first non-empty chunk arrives (using *is_first, echo, chunk_text), so if the
backend returns Complete immediately (e.g., echo=true, max_tokens=0) the echoed
prompt is never emitted; modify the streaming completion path to check the same
conditions (if *is_first && echo) and emit a chunk containing prompt_text before
emitting the terminal/Complete chunk so the echo is always sent even when no
chunks were produced.
- Around line 2388-2390: The loop currently does "if
stopped_indices.contains_key(&index) { continue; }" which skips the
suffix-handling block and therefore drops suffix tokens for choices that also
matched a stop-decoder; change the logic so that stopped_indices is still
checked to prevent streaming further decoded tokens but does NOT short-circuit
the rest of the iteration — remove the continue and instead set/consult a flag
(e.g., is_stopped = stopped_indices.contains_key(&index)) or gate only the
token-emission code, then always run the suffix emission block so suffix bytes
are streamed even when is_stopped is true; apply the same change to the similar
block referenced (lines ~2415-2432) to ensure suffix is never skipped after a
stop match.
🪄 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: 3107ba2e-46c5-4bce-83eb-86c8cf7856b3

📥 Commits

Reviewing files that changed from the base of the PR and between 8ee624e and e0da350.

📒 Files selected for processing (3)
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/response_processing.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: daf9cf6c3f

ℹ️ 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".

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

♻️ Duplicate comments (2)
model_gateway/src/routers/grpc/regular/streaming.rs (2)

2371-2373: ⚠️ Potential issue | 🟠 Major

Don't bypass suffix after a local stop match.

Once an index lands in stopped_indices, this continue skips the suffix block too. Requests that combine stop and suffix still lose the suffix for that choice.

♻️ Suggested fix
-                    if stopped_indices.contains_key(&index) {
-                        continue;
-                    }
+                    let stopped = stopped_indices.contains_key(&index);
 
-                    if let Some(decoder) = stop_decoders.get_mut(&index) {
-                        if let SequenceDecoderOutput::Text(text) = decoder.flush() {
-                            if !text.is_empty() {
-                                let stream_resp = CompletionStreamResponse {
-                                    id: request_id.clone(),
-                                    object: "text_completion".to_string(),
-                                    created,
-                                    choices: vec![CompletionStreamChoice {
-                                        text,
-                                        index,
-                                        logprobs: None,
-                                        finish_reason: None,
-                                    }],
-                                    model: model.clone(),
-                                    system_fingerprint: system_fingerprint.map(String::from),
-                                };
-                                Self::format_completion_sse_into(&mut sse_buffer, &stream_resp);
-                                tx.send(Ok(Bytes::from(sse_buffer.clone())))
-                                    .map_err(|_| "Channel closed".to_string())?;
-                            }
-                        }
-                    }
+                    if !stopped {
+                        if let Some(decoder) = stop_decoders.get_mut(&index) {
+                            if let SequenceDecoderOutput::Text(text) = decoder.flush() {
+                                if !text.is_empty() {
+                                    let stream_resp = CompletionStreamResponse {
+                                        id: request_id.clone(),
+                                        object: "text_completion".to_string(),
+                                        created,
+                                        choices: vec![CompletionStreamChoice {
+                                            text,
+                                            index,
+                                            logprobs: None,
+                                            finish_reason: None,
+                                        }],
+                                        model: model.clone(),
+                                        system_fingerprint: system_fingerprint.map(String::from),
+                                    };
+                                    Self::format_completion_sse_into(&mut sse_buffer, &stream_resp);
+                                    tx.send(Ok(Bytes::from(sse_buffer.clone())))
+                                        .map_err(|_| "Channel closed".to_string())?;
+                                }
+                            }
+                        }
+                    }
 
                     if let Some(sfx) = suffix {
                         let stream_resp = CompletionStreamResponse {
                             id: request_id.clone(),
                             object: "text_completion".to_string(),
@@
                         tx.send(Ok(Bytes::from(sse_buffer.clone())))
                             .map_err(|_| "Channel closed".to_string())?;
                     }
+
+                    if stopped {
+                        continue;
+                    }
 
                     let finish_reason = {

Also applies to: 2398-2448

🤖 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 2371 -
2373, The loop currently uses "if stopped_indices.contains_key(&index) {
continue; }" which skips the subsequent suffix-handling logic and causes
requests with both stop and suffix to lose the suffix; change the control flow
so that when stopped_indices contains index you still run the suffix-handling
block (apply suffix) before skipping further processing — e.g., remove the
immediate continue and instead branch to only skip the later processing while
still executing the suffix code for that index; look for the stopped_indices,
index, and suffix handling in the surrounding loop (also address the same
pattern in the later block around the code referenced 2398-2448) and ensure
suffix is applied even when an index is marked stopped.

2366-2373: ⚠️ Potential issue | 🟠 Major

Emit echo from the Complete path when no chunk was produced.

The echoed prompt is only sent from the chunk branch. If the backend goes straight to Complete—for example echo=true with max_tokens=0—this path emits only the suffix/final events and drops the echoed prompt.

♻️ Suggested fix
                 ProtoResponseVariant::Complete(complete) => {
                     let index = complete.index();
                     total_prompt = total_prompt.max(complete.prompt_tokens());
                     total_completion.record_complete(&complete);
 
+                    let is_first = is_firsts.entry(index).or_insert(true);
+                    if *is_first && echo {
+                        let echo_chunk = CompletionStreamResponse {
+                            id: request_id.clone(),
+                            object: "text_completion".to_string(),
+                            created,
+                            choices: vec![CompletionStreamChoice {
+                                text: prompt_text.to_string(),
+                                index,
+                                logprobs: None,
+                                finish_reason: None,
+                            }],
+                            model: model.clone(),
+                            system_fingerprint: system_fingerprint.map(String::from),
+                        };
+                        Self::format_completion_sse_into(&mut sse_buffer, &echo_chunk);
+                        tx.send(Ok(Bytes::from(sse_buffer.clone())))
+                            .map_err(|_| "Channel closed".to_string())?;
+                        *is_first = false;
+                    }
+
                     if stopped_indices.contains_key(&index) {
                         continue;
                     }
🤖 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/streaming.rs`:
- Around line 2362-2363: The synthetic stop chunk send currently ignores
tx.send's result; change it to treat it like other sends by checking the Result
from tx.send(Ok(Bytes::from(sse_buffer.clone()))) and returning/breaking on Err
so the function exits (allowing upstream to be dropped/aborted) instead of
continuing to EOF and calling mark_completed(); apply the same error-handling
pattern used elsewhere for sends around format_completion_sse_into, final_chunk,
and sse_buffer so disconnects are propagated correctly.

---

Duplicate comments:
In `@model_gateway/src/routers/grpc/regular/streaming.rs`:
- Around line 2371-2373: The loop currently uses "if
stopped_indices.contains_key(&index) { continue; }" which skips the subsequent
suffix-handling logic and causes requests with both stop and suffix to lose the
suffix; change the control flow so that when stopped_indices contains index you
still run the suffix-handling block (apply suffix) before skipping further
processing — e.g., remove the immediate continue and instead branch to only skip
the later processing while still executing the suffix code for that index; look
for the stopped_indices, index, and suffix handling in the surrounding loop
(also address the same pattern in the later block around the code referenced
2398-2448) and ensure suffix is applied even when an index is marked stopped.
🪄 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: 4306c339-e2a5-4950-9e00-b713432c5aef

📥 Commits

Reviewing files that changed from the base of the PR and between e0da350 and daf9cf6.

📒 Files selected for processing (3)
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/response_processing.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 8d2722a1da

ℹ️ 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".

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated
@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Repo admins can enable using credits for code reviews in their settings.

Signed-off-by: VS Chandra Mourya <msrinivasa@together.ai>
@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Repo admins can enable using credits for code reviews in their settings.

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

♻️ Duplicate comments (2)
model_gateway/src/routers/grpc/regular/streaming.rs (2)

2346-2363: ⚠️ Potential issue | 🟠 Major

Keep suffix ahead of the synthetic stop chunk.

When a stop-decoder match happens here, the code sends the terminal finish_reason="stop" chunk immediately and marks the choice as stopped. The Complete arm then short-circuits that index, so requests that combine stop and suffix drop the suffix entirely. Please emit the suffix before this synthetic final chunk, or defer the terminal chunk until the Complete arm so the finish event remains last.

Also applies to: 2371-2373

🤖 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 2346 -
2363, The synthetic terminal chunk with finish_reason="stop" is being emitted
immediately in the stop-decoder match (where stopped_indices.insert(...) and
final_chunk/CompletionStreamChoice are created and sent via
Self::format_completion_sse_into + tx.send), causing any pending suffix to be
dropped by the Complete arm; change the logic so that you either (A) emit any
configured suffix before constructing/sending the synthetic final_chunk (i.e.,
write the suffix into sse_buffer and send it prior to creating the terminal
CompletionStreamResponse), or (B) defer creating/sending the terminal
final_chunk here and let the Complete arm emit the finish chunk last—update the
stop handling around stopped_indices and Self::format_completion_sse_into to
ensure suffix content is always sent before the finish event.

2366-2448: ⚠️ Potential issue | 🟠 Major

Emit the echoed prompt when the backend goes straight to Complete.

The prompt is only prepended in the Chunk arm. If a backend emits no chunks at all—for example, echo=true with max_tokens=0—this path returns only the terminal chunk (and maybe suffix) without ever streaming the prompt. If this index is still first-seen in the Complete arm, send a prompt chunk before suffix/finalization.

🤖 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/streaming.rs`:
- Around line 2275-2276: The code reads completion_request.logprobs but never
includes per-token logprobs in the streamed choices (CompletionStreamChoice),
causing stream=true with logprobs set to return incomplete payloads; either
implement plumbing of per-chunk logprobs through the streaming path or
explicitly reject the unsupported combination. Fix by adding an early check
where the stream handler inspects completion_request.logprobs (the same spot at
the referenced occurrences around the other blocks at 2331-2335 and 2437-2441)
and if logprobs.is_some() return a clear error (e.g., unsupported parameter
combination) or, if you choose to implement, propagate the backend’s per-token
logprobs into each emitted CompletionStreamChoice.logprobs instead of leaving it
None. Ensure the check/propagation happens in the streaming response
construction that creates CompletionStreamChoice so the behavior is consistent
across all referenced blocks.

---

Duplicate comments:
In `@model_gateway/src/routers/grpc/regular/streaming.rs`:
- Around line 2346-2363: The synthetic terminal chunk with finish_reason="stop"
is being emitted immediately in the stop-decoder match (where
stopped_indices.insert(...) and final_chunk/CompletionStreamChoice are created
and sent via Self::format_completion_sse_into + tx.send), causing any pending
suffix to be dropped by the Complete arm; change the logic so that you either
(A) emit any configured suffix before constructing/sending the synthetic
final_chunk (i.e., write the suffix into sse_buffer and send it prior to
creating the terminal CompletionStreamResponse), or (B) defer creating/sending
the terminal final_chunk here and let the Complete arm emit the finish chunk
last—update the stop handling around stopped_indices and
Self::format_completion_sse_into to ensure suffix content is always sent before
the finish event.
🪄 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: 2f671285-2c84-4c08-b891-f5004d95fbb7

📥 Commits

Reviewing files that changed from the base of the PR and between 8d2722a and 39bb602.

📒 Files selected for processing (3)
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/response_processing.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
…rors in streaming

Signed-off-by: VS Chandra Mourya <msrinivasa@together.ai>
@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Repo admins can enable using credits for code reviews in their settings.

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs Outdated
Comment thread model_gateway/src/routers/grpc/regular/streaming.rs
Comment on lines +2343 to +2356
let stream_resp = CompletionStreamResponse {
id: request_id.clone(),
object: "text_completion".to_string(),
created,
choices: vec![CompletionStreamChoice {
text: std::mem::take(&mut chunk_text),
index,
logprobs: None,
finish_reason: None,
}],
model: model.clone(),
system_fingerprint: system_fingerprint.map(String::from),
usage: None,
};

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: Per-token String allocations in the streaming hot path. Every SSE chunk constructs a new CompletionStreamResponse with request_id.clone(), model.clone(), "text_completion".to_string(), and system_fingerprint.map(String::from). These are new heap allocations on every token.

Consider pre-allocating these once before the loop (e.g., as owned Strings) and using Cow::Borrowed or simply serializing directly into the SSE buffer with a custom writer that avoids the intermediate struct allocation. At minimum, "text_completion" could be a &'static str if CompletionStreamResponse used Cow<'a, str> for the object field.

Comment on lines +2534 to +2543
Metrics::record_streaming_metrics(StreamingMetricsParams {
router_type: metrics_labels::ROUTER_GRPC,
backend_type: self.backend_type,
model_id: model,
endpoint: metrics_labels::ENDPOINT_COMPLETIONS,
ttft: first_token_time.map(|t| t.duration_since(start_time)),
generation_duration: start_time.elapsed(),
input_tokens: Some(total_prompt as u64),
output_tokens: total_completion.total() as u64,
});

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit: Metrics are recorded unconditionally even when the stream errored out mid-way (returned Err before reaching this point). Since ? short-circuits on stream/channel errors earlier in the function, this code only runs on the success path — which is correct. However, first_token_time could be None if the stream completed with only Complete messages (no Chunks), in which case TTFT is reported as None. Worth confirming that the metrics pipeline handles a None TTFT gracefully for completions (it likely does, since embeddings would have the same shape).

Comment thread model_gateway/src/routers/grpc/regular/streaming.rs

result
}

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We should plan to split this file as a follow up. A lot of functions here can be shared. cc @slin1237

@CatherineSue CatherineSue left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the PR. Will take a detailed look shortly.

@CatherineSue CatherineSue left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Overall LGTM. Only have one comment regarding the max_tokens=0. Please take a look.

As a follow up, we can add completions to e2e-tests in the next PR. Thanks for the contribution!

@github-actions github-actions Bot added the protocols Protocols crate changes label Apr 1, 2026
@vschandramourya

Copy link
Copy Markdown
Collaborator Author

@CatherineSue added echo handling in the Complete arm for the max_tokens=0 case. Sure E2E will be a follow up PR I can raise. Thank you for the review.

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 735db19778

ℹ️ 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".

Comment on lines +2292 to +2294
if stopped_indices.contains(&index) {
continue;
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Keep counting tokens for stopped vLLM choices

This early continue skips all later chunks for an index once local stop matching fires, which means total_completion.record_chunk(&chunk) no longer runs for that index. In this code path, usage and metrics rely on CompletionTokenTracker, and for vLLM that tracker is chunk-accumulated rather than overwritten by Complete; so when backend generation continues after local stop (the nearby comment already notes Complete may later report length), the final include_usage chunk under-reports completion tokens. This affects streaming requests with stop sequences on vLLM, especially when stream_options.include_usage=true.

Useful? React with 👍 / 👎.

@vschandramourya
vschandramourya force-pushed the mourya/cmp-5 branch 2 times, most recently from 28e02f7 to e0e82c9 Compare April 1, 2026 19:57

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: e0e82c9695

ℹ️ 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".

tx.send(Ok(Bytes::from(sse_buffer.clone())))
.map_err(|_| "Channel closed".to_string())?;
}
_ => continue,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Propagate backend error frames during completion streaming

The wildcard arm here also swallows ProtoResponseVariant::Error, so if the worker sends an in-band error during decode, this loop silently ignores it and the task still emits a terminal [DONE]. In that scenario clients receive a seemingly successful but truncated stream instead of an error event, which diverges from the other streaming handlers in this file that return Err(...) on Error variants.

Useful? React with 👍 / 👎.

Comment on lines +2588 to +2589
ProtoResponseVariant::Complete(_) => break,
_ => continue,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Fail dual completion stream when prefill returns error

In PD mode, the prefill loop only exits on Complete and treats every other variant as continue, which includes ProtoResponseVariant::Error. If prefill fails in-band, the function currently ignores that failure and proceeds to decode, potentially returning a normal-looking stream for a failed request; this should return an error immediately like the dual messages path does.

Useful? React with 👍 / 👎.

Signed-off-by: VS Chandra Mourya <msrinivasa@together.ai>

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: f2aedeb0d6

ℹ️ 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".

Comment on lines +2476 to +2480
"stop" | "length" | "content_filter" => t.to_string(),
other => {
warn!(unexpected_finish_reason = other, "Unmapped finish_reason type from backend, defaulting to stop");
"stop".to_string()
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Preserve unknown finish reasons in completion streams

When the backend returns a finish reason outside stop|length|content_filter (for example a new backend-specific value), this branch rewrites it to "stop" instead of forwarding the original signal. That makes streamed /v1/completions responses report a successful stop even when the backend indicated a different termination reason, and it diverges from non-streaming completion handling in processor.rs which preserves unknown reasons. Forwarding unknown values (or mapping to an explicit fallback) avoids silently misclassifying termination state for clients.

Useful? React with 👍 / 👎.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

grpc gRPC client and router changes model-gateway Model gateway crate changes protocols Protocols crate changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants