Skip to content

refactor(grpc): deduplicate Harmony Chat Completion streaming decode logic - #594

Merged
CatherineSue merged 2 commits into
mainfrom
chang/refactor-streaming-dedup
Mar 3, 2026
Merged

CatherineSue merged 2 commits into
mainfrom
chang/refactor-streaming-dedup

Conversation

@CatherineSue

@CatherineSue CatherineSue commented Mar 3, 2026 •

Copy link
Copy Markdown
Member

Description

Problem

PR #592 review (Gemini, high priority) flagged significant code duplication between process_single_stream and process_dual_stream in streaming.rs. The decode loop, chunk parsing, delta emission, completion handling, finalization, usage emission, and metrics recording were nearly identical (~160 lines duplicated). The only differences:

  1. dual_stream has a prefill phase that pre-populates prompt_tokens and cached_tokens
  2. single_stream populates prompt_tokens/cached_tokens from Complete messages; dual_stream doesn't
  3. dual_stream marks two streams complete (decode first, then prefill); single_stream marks one

Solution

Extract the shared decode logic into process_chat_decode_stream(). Callers only set up prompt_tokens/cached_tokens and handle prefill-specific cleanup.

The prompt_tokens/cached_tokens difference is handled via entry().or_insert_with():

  • single_stream: maps start empty → always inserts (same as original .insert())
  • dual_stream: maps pre-populated from prefill → or_insert_with is a no-op (closure never called)

Changes

  • Added process_chat_decode_stream() — shared helper containing the full decode loop, metrics, and usage emission
  • Simplified process_single_stream() — now just creates empty maps and delegates to the helper
  • Simplified process_dual_stream() — keeps prefill phase, delegates decode to the helper, marks prefill stream completed after

Performance note: The dual-stream path now does an entry().or_insert_with() lookup in the Complete arm where the original code didn't touch those maps at all. This is one u32 hash + comparison per Complete message per map — sub-nanosecond and completely negligible compared to any network I/O in the stream. The closure is lazy so complete_wrapper.prompt_tokens() / cached_tokens() are never called when the key already exists.

Test Plan

  • cargo build -p smg — passes
  • cargo clippy -p smg -- -D warnings — passes (zero warnings)
  • cargo test -p smg — all tests pass
  • pre-commit run --all-files — all hooks pass

Note: PR #592 will need to rebase on this after merge.

Checklist
  • cargo +nightly fmt passes
  • cargo clippy --all-targets --all-features -- -D warnings passes
  • (Optional) Documentation updated

Summary by CodeRabbit

  • Refactor
    • Consolidated streaming decode logic into a shared handler, reducing duplication and centralizing per-chunk parsing, delta emission, finalization, and usage aggregation.
    • Synchronized token-tracking across decode and completion phases and clarified stream completion semantics to improve streaming reliability and preserve existing per-chunk behavior.

@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.

@github-actions github-actions Bot added grpc gRPC client and router changes model-gateway Model gateway crate changes labels Mar 3, 2026
@gemini-code-assist

Copy link
Copy Markdown
Contributor

Summary of Changes

Hello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed!

This pull request significantly refactors the gRPC Chat Completion streaming logic by addressing and eliminating substantial code duplication. The core decode loop, including metrics and usage emission, has been extracted into a new, shared helper function. This change enhances the maintainability and clarity of the streaming implementation across both single and dual stream processing paths, making the codebase more robust and easier to understand.

Highlights

  • Code Duplication Eliminated: Extracted the shared Chat Completion decode logic from process_single_stream() and process_dual_stream() into a new helper function, process_chat_decode_stream().
  • Simplified Single Stream Processing: The process_single_stream() function was simplified to initialize empty token maps and then delegate its core streaming decode functionality to the new shared helper.
  • Simplified Dual Stream Processing: The process_dual_stream() function now handles its prefill phase and then delegates the subsequent decode phase to the new shared helper, ensuring proper stream completion.
  • Optimized Token Map Handling: Updated the handling of prompt_tokens and cached_tokens within the shared decode stream to use entry().or_insert_with(), which efficiently populates values only if not already present, respecting pre-populated values from the dual stream's prefill phase.
Changelog
  • model_gateway/src/routers/grpc/harmony/streaming.rs
    • Introduced process_chat_decode_stream to encapsulate the common streaming decode logic, including chunk parsing, delta emission, completion handling, usage emission, and metrics recording.
    • Modified process_single_stream to initialize empty token maps and then call process_chat_decode_stream for its core functionality.
    • Updated process_dual_stream to handle its prefill phase, then delegate the decode phase to process_chat_decode_stream, and finally mark the prefill stream as completed.
    • Adjusted how prompt_tokens and cached_tokens are populated within the shared decode stream, using entry().or_insert_with() to respect pre-populated values from the dual stream's prefill phase.
    • Changed parameter passing for dispatch and original_request in calls to emit_chunk_delta and emit_final_chunk within the new shared function to use direct references instead of double references.
Activity
Using Gemini Code Assist

The 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 /gemini <command> or @gemini-code-assist <command>. Below is a summary of the supported commands on the current page.

Feature Command Description
Code Review /gemini review Performs a code review for the current pull request in its current state.
Pull Request Summary /gemini summary Provides a summary of the current pull request in its current state.
Comment @gemini-code-assist Responds in comments when explicitly tagged, both in pull request comments and review comments.
Help /gemini help Displays a list of available commands.

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 .gemini/ folder in the base of the repository. Detailed instructions can be found here.

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

  1. Review the Privacy Notices, Generative AI Prohibited Use Policy, Terms of Service, and learn how to configure Gemini Code Assist in GitHub here. Gemini can make mistakes, so double check it and use code with caution. ↩

@coderabbitai

coderabbitai Bot commented Mar 3, 2026 •

Copy link
Copy Markdown

Caution

Review failed

The pull request is closed.

ℹ️ Recent review info

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between 770ff9f and 61eacad.

📒 Files selected for processing (1)
  • model_gateway/src/routers/grpc/harmony/streaming.rs

📝 Walkthrough

Walkthrough

Introduces process_chat_decode_stream to centralize per-chunk decode/parsing, delta emission, finalization, and usage aggregation; updates process_single_stream and process_dual_stream signatures to delegate decoding, harmonizes prompt/cached token handling, and adjusts stream completion points.

Changes

Cohort / File(s) Summary
Streaming processor refactor
model_gateway/src/routers/grpc/harmony/streaming.rs
Added process_chat_decode_stream and moved per-chunk parsing/delta emission/finalization/usage aggregation into it. Updated process_single_stream and process_dual_stream signatures to delegate decode handling, unified prompt_tokens/cached_tokens population, and adjusted which stream is marked completed in each flow.
Manifest note
Cargo.toml
Referenced previously; verify any dependency/version edits if present.

Sequence Diagram(s)

sequenceDiagram
    participant Client
    participant HarmonyProcessor
    participant PrefillStream
    participant DecodeStream
    participant Parser as HarmonyParserAdapter
    participant TX as mpsc::UnboundedSender

    Client->>HarmonyProcessor: start single or dual streaming
    alt dual flow
        HarmonyProcessor->>PrefillStream: consume prefill until ready
    end
    HarmonyProcessor->>DecodeStream: read decode chunks
    DecodeStream-->>HarmonyProcessor: chunk bytes
    HarmonyProcessor->>Parser: parse chunk -> deltas / final / usage
    Parser-->>HarmonyProcessor: delta events, final event, usage
    HarmonyProcessor->>TX: emit delta bytes
    HarmonyProcessor->>DecodeStream: mark_completed (when decode done)
    alt dual flow post-decode
        HarmonyProcessor->>PrefillStream: mark_completed
    end
    HarmonyProcessor->>TX: emit final bytes and usage
Loading

Estimated code review effort

🎯 3 (Moderate) | ⏱️ ~20 minutes

Possibly related PRs

Suggested reviewers

  • key4ng
  • slin1237

Poem

🐰 I hopped through bytes and parsed each stream,
Gathered tiny deltas into one bright beam,
Tokens counted, endings sung with care,
One helper hums — less juggling to bear,
Carrots for CI, and tests if you dare! 🥕

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Title check ✅ Passed The title directly and concisely summarizes the main change: refactoring to deduplicate streaming decode logic in gRPC Harmony Chat Completion handling.

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

✨ Finishing Touches
  • 📝 Generate docstrings (stacked PR)
  • 📝 Generate docstrings (commit on current branch)
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Post copyable unit tests in a comment
  • Commit unit tests in branch chang/refactor-streaming-dedup

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

@mergify

This comment was marked as resolved.

@CatherineSue
CatherineSue force-pushed the chang/refactor-streaming-dedup branch from f575849 to a93b339 Compare March 3, 2026 20:17

@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 refactors the chat completion streaming logic, successfully deduplicating code by introducing the process_chat_decode_stream helper function. No security vulnerabilities were found. The solution is clean and elegant, with a suggestion to improve the maintainability of the new helper function by encapsulating its state into a dedicated struct.

…ess_chat_decode_stream

The decode loop, chunk parsing, delta emission, completion handling,
finalization, usage emission, and metrics recording were nearly identical
(~160 lines duplicated) between process_single_stream and process_dual_stream.

Extract the shared decode logic into process_chat_decode_stream(). The only
meaningful difference — prompt_tokens/cached_tokens population — is handled
via entry().or_insert_with(): single-stream passes empty maps (always inserts),
dual-stream passes pre-populated maps from prefill (no-op on Complete).

Signed-off-by: Chang Su <chang.s.su@oracle.com>
@CatherineSue
CatherineSue force-pushed the chang/refactor-streaming-dedup branch from a93b339 to 770ff9f Compare March 3, 2026 20:39

@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.

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
model_gateway/src/routers/grpc/harmony/streaming.rs (1)

205-206: 🧹 Nitpick | 🔵 Trivial

Remove unused finish_reasons map from the shared decode helper.

finish_reasons is populated (Line 275) but never read, so it adds dead state in a hot streaming path.

♻️ Suggested cleanup
-        let mut finish_reasons: HashMap<u32, Option<String>> = HashMap::new();
@@
-                    finish_reasons
-                        .insert(index, Some(complete_wrapper.finish_reason().to_string()));

Also applies to: 275-277

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@model_gateway/src/routers/grpc/harmony/streaming.rs` around lines 205 - 206,
The variable finish_reasons in the shared decode helper is unused and introduces
dead state; remove its declaration (the HashMap<u32, Option<String>>) and any
code that inserts into or references finish_reasons (the population logic around
where finish_reasons is set) so the streaming path only keeps matched_stops
(HashMap<u32, Option<serde_json::Value>>) and related logic; update any function
signatures or variables that referenced finish_reasons to stop expecting it and
run tests/build to ensure no remaining references.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Outside diff comments:
In `@model_gateway/src/routers/grpc/harmony/streaming.rs`:
- Around line 205-206: The variable finish_reasons in the shared decode helper
is unused and introduces dead state; remove its declaration (the HashMap<u32,
Option<String>>) and any code that inserts into or references finish_reasons
(the population logic around where finish_reasons is set) so the streaming path
only keeps matched_stops (HashMap<u32, Option<serde_json::Value>>) and related
logic; update any function signatures or variables that referenced
finish_reasons to stop expecting it and run tests/build to ensure no remaining
references.

ℹ️ Review info

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

📥 Commits

Reviewing files that changed from the base of the PR and between a93b339 and 770ff9f.

📒 Files selected for processing (1)
  • model_gateway/src/routers/grpc/harmony/streaming.rs

@CatherineSue CatherineSue changed the title refactor(grpc): deduplicate Chat Completion streaming decode logic refactor(grpc): deduplicate Harmony Chat Completion streaming decode logic Mar 3, 2026
finish_reasons was populated but never read back; the finish reason is
consumed directly from complete_wrapper.finish_reason() instead.

Signed-off-by: Chang Su <chang.s.su@oracle.com>
@CatherineSue
CatherineSue merged commit d50e4d8 into main Mar 3, 2026
9 of 11 checks passed
@CatherineSue
CatherineSue deleted the chang/refactor-streaming-dedup branch March 3, 2026 22:03
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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant