Skip to content

feat(multimodal): add EPD encode routing - #1852

Merged
slin1237 merged 9 commits into
smg-project:mainfrom
chenht2022:epd-routing-clean
Jul 7, 2026
Merged

slin1237 merged 9 commits into
smg-project:mainfrom
chenht2022:epd-routing-clean

Conversation

@chenht2022

@chenht2022 chenht2022 commented Jun 28, 2026 •

Copy link
Copy Markdown
Member

High-level approach

This PR adds the SMG gateway side of TokenSpeed Encode-Prefill-Decode (EPD) routing for multimodal requests.

For TokenSpeed EPD, SMG owns the encode routing decision: the gateway routes multimodal encode work to encode workers and sends the paired prefill/decode request with rendezvous metadata; the prefill worker waits for the encoded payload through the TokenSpeed/Mooncake transfer path.

At a high level, the request flow is:

  1. Request building detects EPD mode and builds an ExecutionPlan::EncodePrefillDecode for supported multimodal chat/messages requests.
  2. Worker selection chooses encode workers per multimodal item, then chooses compatible prefill/decode workers.
  3. Request execution dispatches TokenSpeed encode RPCs, wires the encode-to-prefill rendezvous metadata into the prefill request, strips raw pixel payloads before prefill, and then continues through the existing prefill/decode execution path.

Encode assignment is per multimodal item rather than per request. The default encode policy is consistent hashing, so repeated item hashes keep affinity with the same encode worker when possible.

Interfaces changed and added

This PR adds EPD as a new routing shape alongside the existing prefill/decode routing path.

Key shared gateway interfaces:

  • Adds RoutingMode::EncodePrefillDecode, parallel to RoutingMode::PrefillDecode, to represent a three-role encode/prefill/decode topology.
  • Adds WorkerType::Encode, which extends shared worker discovery, worker registry stats, health/readiness checks, metrics labels, and job queue expansion.
  • Adds ExecutionPlan / ExecutionPlanKind as the request handoff between request-building and execution stages, so the common gRPC pipeline can represent single-worker, prefill/decode, and encode/prefill/decode flows.
  • Extends policy handling with an encode policy in PolicyRegistry; encode defaults to consistent_hashing.
  • Extends config/CLI/Python bindings with EPD fields, including encode URLs/selectors and encode policy.

TokenSpeed-specific interfaces:

  • Adds tokenspeed_encoder.proto and generated Rust/Python exports for the TokenSpeedEncoder service.
  • Adds TokenSpeedEncoderClient for gateway-side encode dispatch.
  • Extends TokenSpeed scheduler/bootstrap metadata so encode workers and prefill workers can rendezvous through the disaggregated transfer path.
  • Adds a Python TokenSpeedEncoderServicer and mounts it when TokenSpeed starts in encode mode.

How to use

Start TokenSpeed workers in encode/prefill/decode roles:

# Encode worker: vision tower only
python3 -m smg_grpc_servicer.tokenspeed ... \
  --disaggregation-mode encode \
  --disaggregation-bootstrap-port 18995

# Prefill worker
python3 -m smg_grpc_servicer.tokenspeed ... \
  --disaggregation-mode prefill \
  --disaggregation-bootstrap-port 19311

# Decode worker
python3 -m smg_grpc_servicer.tokenspeed ... \
  --disaggregation-mode decode

Launch SMG in EPD mode:

python3 -m smg launch \
  --epd-disaggregation \
  --encode grpc://ENCODE_HOST:50104 18995 \
  --prefill grpc://PREFILL_HOST:50101 19311 \
  --decode grpc://DECODE_HOST:50111 \
  --model-path "$MODEL_PATH" \
  --tokenizer-path "$MODEL_PATH" \
  --host 0.0.0.0 \
  --port 8000

Mooncake is the default disaggregation transfer backend for this EPD path.

Test results

End-to-end EPD was tested with the paired SMG EPD gateway/servicer stack and the TokenSpeed runtime changes in lightseekorg/tokenspeed#548.

Representative setup:

  • Hardware: NVIDIA B200 GPUs
  • Model: Qwen3.5-122B-A10B-NVFP4
  • Load: unique 1080p images, 8 images/request, max_tokens=1
  • Transport: encode-to-prefill embeddings over Mooncake

Scaling results:

Topology Peak throughput
1E2P2D 5.37 req/s
2E4P2D 10.63 req/s
3E6P2D 15.61 req/s
4E8P2D 18.91 req/s

Potential follow-ups

  • feat(multimodal): optimize EPD encode routing #1853 follows up on EPD routing performance, including the optional RDMA/NIXL path for gateway-to-encode pixel transfer.
  • This SMG PR and feat(disaggregation): add EPD encode pipeline lightseekorg/tokenspeed#548 are coupled because they change the cross-repo SMG/TokenSpeed interface. The TokenSpeed E2E CI in this SMG PR is expected to fail because SMG CI still tests against the old TokenSpeed side. The landing plan is to land this SMG PR first with the known TokenSpeed E2E CI failure, then publish/bump tokenspeed-smg* while landing TokenSpeed PR, and finally bump SMG's TOKENSPEED_REF.

Summary by CodeRabbit

  • New Features
    • Added Encode-Prefill-Decode (EPD) disaggregated mode with dedicated encode workers, gRPC EPD routing, and updated readiness/health behavior.
    • Introduced a new gRPC TokenSpeed encoder service and end-to-end encode request flow for multimodal inputs.
    • Added new configuration/Python binding options for EPD (encode URLs, selectors, and per-stage encode policy).
  • Bug Fixes
    • Improved multimodal request assembly and bootstrap rendezvous handling so encode/prefill/decode stages share the correct metadata and rendezvous parameters.

@coderabbitai

coderabbitai Bot commented Jun 28, 2026 •

Copy link
Copy Markdown

Warning

Review limit reached

@slin1237, you've reached your PR review limit, so we couldn't start this review.

Next review available in: 1 minute

Enable usage-based reviews in Billing to review now. Otherwise, wait until the next included review is available.
You're only billed for reviews past your plan's rate limits ($0.25/file).

How can I continue?

After more reviews become available, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews.

How do review limits work?

CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability.

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window.

Please refer docs for additional details.

Review details
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 18179629-c914-4374-924e-0c3239ff73e8

📥 Commits

Reviewing files that changed from the base of the PR and between 83db4b2 and cfecdc8.

📒 Files selected for processing (50)
  • bindings/python/src/lib.rs
  • crates/grpc_client/build.rs
  • crates/grpc_client/proto/tokenspeed_encoder.proto
  • crates/grpc_client/proto/tokenspeed_scheduler.proto
  • crates/grpc_client/python/smg_grpc_proto/__init__.py
  • crates/grpc_client/src/lib.rs
  • crates/grpc_client/src/tokenspeed_encoder.rs
  • crates/grpc_client/src/tokenspeed_scheduler.rs
  • crates/protocols/src/worker.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/server.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
  • model_gateway/src/config/types.rs
  • model_gateway/src/config/validation.rs
  • model_gateway/src/health.rs
  • model_gateway/src/main.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/policies/registry.rs
  • model_gateway/src/routers/factory.rs
  • model_gateway/src/routers/grpc/common/response_collection.rs
  • model_gateway/src/routers/grpc/common/stages/client_acquisition.rs
  • model_gateway/src/routers/grpc/common/stages/dispatch_metadata.rs
  • model_gateway/src/routers/grpc/common/stages/helpers.rs
  • model_gateway/src/routers/grpc/common/stages/mod.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs
  • model_gateway/src/routers/grpc/context.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/harmony/stages/request_building.rs
  • model_gateway/src/routers/grpc/harmony/streaming.rs
  • model_gateway/src/routers/grpc/mod.rs
  • model_gateway/src/routers/grpc/multimodal.rs
  • model_gateway/src/routers/grpc/pd_router.rs
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/proto_wrapper.rs
  • model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/classify/response_processing.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/embedding/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/generate/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/request_building.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs
  • model_gateway/src/routers/router_manager.rs
  • model_gateway/src/service_discovery.rs
  • model_gateway/src/worker/manager.rs
  • model_gateway/src/worker/registry.rs
  • model_gateway/src/worker/service.rs
  • model_gateway/src/worker/worker.rs
  • model_gateway/src/workflow/job_queue.rs
📝 Walkthrough

Walkthrough

Adds encode-prefill-decode routing, execution-plan plumbing, TokenSpeed encoder servicing, and related config, discovery, multimodal, and streaming updates.

Changes

EPD config, discovery, and routing

Layer / File(s) Summary
Config, policy, discovery, and routing support
bindings/python/src/lib.rs, model_gateway/src/config/types.rs, model_gateway/src/config/validation.rs, model_gateway/src/main.rs, model_gateway/src/service_discovery.rs, model_gateway/src/policies/registry.rs, model_gateway/src/routers/factory.rs, model_gateway/src/routers/router_manager.rs, model_gateway/src/workflow/job_queue.rs, model_gateway/src/health.rs, model_gateway/src/observability/metrics.rs, model_gateway/src/worker/{manager.rs,registry.rs,service.rs,worker.rs}
Adds EPD routing mode, encode selectors and policy plumbing, disaggregated service discovery, encode worker accounting, router selection, and startup/binding configuration.

Execution planning and request flow

Layer / File(s) Summary
Execution plan and worker-selection model
model_gateway/src/routers/grpc/{context.rs,proto_wrapper.rs,epd_encode.rs}, model_gateway/src/routers/grpc/common/stages/{helpers.rs,client_acquisition.rs,dispatch_metadata.rs,mod.rs,request_execution.rs,worker_selection.rs}
Replaces raw proto handoff with execution plans, adds encode-dispatch planning, and renames dual/disaggregated selection paths.
Request building, multimodal assembly, and pipelines
model_gateway/src/routers/grpc/{multimodal.rs,pd_router.rs,pipeline.rs}, model_gateway/src/routers/grpc/{regular,harmony}/stages/*
Updates request builders and pipelines to emit execution plans, choose encode-aware multimodal assembly, inject rendezvous metadata, and construct EPD routers.
Streaming and response collection
model_gateway/src/routers/grpc/{common/response_collection.rs,regular/streaming.rs,harmony/streaming.rs}
Renames prefill/decode streaming variants and routes response collection through PrefillDecode.

TokenSpeed encoder integration

Layer / File(s) Summary
Encoder proto, client, servicer, and startup wiring
crates/grpc_client/{build.rs,proto/tokenspeed_{encoder,scheduler}.proto,python/smg_grpc_proto/__init__.py,src/{lib.rs,tokenspeed_encoder.rs,tokenspeed_scheduler.rs}}, grpc_servicer/smg_grpc_servicer/tokenspeed/{encoder_servicer.py,server.py,servicer.py}
Adds the encoder RPC contract and client, exposes it in generated modules, serves it in encode mode, and threads encode/bootstrap fields through TokenSpeed request handling.

Estimated code review effort: 5 (Critical) | ~120 minutes

Possibly related issues

Possibly related PRs

Suggested labels: multimodal, model-gateway

Suggested reviewers: CatherineSue, key4ng, slin1237, gongwei-130

Poem

I hopped through plans with ears held high,
Encode and decode took to the sky.
A bunny stitched rooms, routes, and light,
And EPD hummed through the night. 🐇

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title is concise and accurately captures the main change: adding EPD encode routing for multimodal requests.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

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.

@github-actions github-actions Bot added grpc gRPC client and router changes protocols Protocols crate changes model-gateway Model gateway crate changes labels Jun 28, 2026

@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 introduces support for the Encode-Prefill-Decode (EPD) disaggregated mode, allowing vision-tower-only encode workers to run the vision tower and ship image embeddings to prefill workers over Mooncake. It adds the TokenSpeedEncoder gRPC service, a Python servicer, a Rust client, and updates the gateway's routing, validation, and service discovery to support EPD. The review feedback highlights several critical improvements: awaiting the asynchronous ZMQ send in the Python servicer to prevent unawaited coroutine warnings, using an RAII guard in epd_encode.rs to prevent shared memory leaks upon cancellation or panic, adopting tokio::task::JoinSet instead of join_all for concurrent task spawning, and establishing channel connections concurrently rather than sequentially to reduce latency.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py Outdated
Comment on lines +92 to +128
pub(crate) async fn dispatch(
mut self,
endpoint: String,
bootstrap_room: i64,
) -> std::result::Result<(), String> {
match &mut self {
Self::TokenSpeed {
item,
shm_enabled,
cleanup_on_drop,
} => {
let item = item
.take()
.ok_or_else(|| "encode item was already dispatched".to_string())?;
*cleanup_on_drop = false;
let request = tokenspeed_encoder::EncodeRequest {
request_id: format!("encode-{}", Uuid::now_v7()),
mm_inputs: Some(
TokenSpeedMultimodalData {
items: vec![item],
shm_enabled: *shm_enabled,
}
.into_proto(),
),
items: vec![tokenspeed_encoder::EncodeItemAssignment { bootstrap_room }],
};
let shm_handles = request
.mm_inputs
.as_ref()
.map(collect_tokenspeed_multimodal_inputs_shm_handles)
.unwrap_or_default();
let result = send_tokenspeed_encode_rpc(endpoint, request).await;
cleanup_tokenspeed_shm_handles(&shm_handles);
result
}
}
}

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.

high

If the send_tokenspeed_encode_rpc future is cancelled (e.g., due to client disconnect) or panics, the manual cleanup of shared memory handles on line 124 will never be reached, leading to a resource leak. Using a local RAII guard ensures that cleanup_tokenspeed_shm_handles is always executed when the scope is exited, regardless of cancellation or panics.

    pub(crate) async fn dispatch(
        mut self,
        endpoint: String,
        bootstrap_room: i64,
    ) -> std::result::Result<(), String> {
        match &mut self {
            Self::TokenSpeed {
                item,
                shm_enabled,
                cleanup_on_drop,
            } => {
                let item = item
                    .take()
                    .ok_or_else(|| "encode item was already dispatched".to_string())?;
                *cleanup_on_drop = false;
                let request = tokenspeed_encoder::EncodeRequest {
                    request_id: format!("encode-{}", Uuid::now_v7()),
                    mm_inputs: Some(
                        TokenSpeedMultimodalData {
                            items: vec![item],
                            shm_enabled: *shm_enabled,
                        }
                        .into_proto(),
                    ),
                    items: vec![tokenspeed_encoder::EncodeItemAssignment { bootstrap_room }],
                };
                let shm_handles = request
                    .mm_inputs
                    .as_ref()
                    .map(collect_tokenspeed_multimodal_inputs_shm_handles)
                    .unwrap_or_default();

                struct ShmGuard<T>(T);
                impl<T> Drop for ShmGuard<T> {
                    fn drop(&mut self) {
                        cleanup_tokenspeed_shm_handles(&self.0);
                    }
                }
                let _guard = ShmGuard(shm_handles);

                send_tokenspeed_encode_rpc(endpoint, request).await
            }
        }
    }

Comment thread model_gateway/src/routers/grpc/common/stages/request_execution.rs Outdated
Comment on lines +87 to +90
let mut built = Vec::with_capacity(ENCODE_CONNS_PER_ENDPOINT);
for _ in 0..ENCODE_CONNS_PER_ENDPOINT {
built.push(crate::channel::connect_channel(endpoint).await?);
}

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.

medium

Connecting to the channels sequentially in a loop introduces unnecessary latency on the connection path. Connecting concurrently using futures::future::try_join_all will significantly speed up connection establishment.

        let mut futures = Vec::with_capacity(ENCODE_CONNS_PER_ENDPOINT);
        for _ in 0..ENCODE_CONNS_PER_ENDPOINT {
            futures.push(crate::channel::connect_channel(endpoint));
        }
        let built = futures::future::try_join_all(futures).await?;

@chenht2022
chenht2022 marked this pull request as ready for review June 28, 2026 14:36
@chenht2022
chenht2022 requested a review from CatherineSue as a code owner June 28, 2026 14:36
Copilot AI review requested due to automatic review settings June 28, 2026 14:36

Copilot AI 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.

Pull request overview

This PR adds an Encode–Prefill–Decode (EPD) disaggregation routing mode for the TokenSpeed gRPC backend. It extends the gateway control-plane (worker discovery/selection, routing, execution planning) and the TokenSpeed runtime/proto surface so multimodal requests can be encoded on dedicated vision-only workers before the existing prefill/decode path runs.

Changes:

  • Introduces RoutingMode::EncodePrefillDecode and WorkerType::Encode, including selection, stats, and health gating.
  • Refactors gRPC request building/execution around an ExecutionPlan to support Single/PD/EPD flows, including encode dispatch + KV rendezvous wiring for TokenSpeed.
  • Adds TokenSpeed EPD encode RPC/proto/client and a Python encode servicer to run vision-tower encode and ship embeddings via Mooncake.

Reviewed changes

Copilot reviewed 49 out of 49 changed files in this pull request and generated 4 comments.

Show a summary per file
File Description
model_gateway/src/workflow/job_queue.rs Adds EPD worker URL expansion and registers encode worker type into specs.
model_gateway/src/worker/worker.rs Adds metrics label mapping for WorkerType::Encode.
model_gateway/src/worker/service.rs Adjusts worker listing/counting behavior to account for encode workers.
model_gateway/src/worker/registry.rs Tracks encode worker counts in registry stats and adds coverage test.
model_gateway/src/worker/manager.rs Propagates encode worker type into worker manager payloads.
model_gateway/src/service_discovery.rs Generalizes PD discovery to “disaggregated” and adds encode pod type/selector + bootstrap-port handling.
model_gateway/src/routers/router_manager.rs Adds GRPC_EPD router id and weighted selection logic for EPD readiness.
model_gateway/src/routers/grpc/regular/streaming.rs Renames “Dual” streaming terminology to “PrefillDecode” for clarity with EPD.
model_gateway/src/routers/grpc/regular/stages/request_building.rs Threads ExecutionPlanKind through chat/generate request-building stage creation.
model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs Adds EPD encode planning + bootstrap/KV rendezvous injection for Messages API; emits an ExecutionPlan.
model_gateway/src/routers/grpc/regular/stages/generate/request_building.rs Moves generate request building to emit an ExecutionPlan and injects TokenSpeed KV rendezvous when applicable.
model_gateway/src/routers/grpc/regular/stages/embedding/request_building.rs Moves embedding request building to ExecutionPlan::embed.
model_gateway/src/routers/grpc/regular/stages/completion/request_building.rs Moves completion request building to ExecutionPlan and injects TokenSpeed KV rendezvous for EPD deployments.
model_gateway/src/routers/grpc/regular/stages/classify/response_processing.rs Updates worker-selection variant naming (Dual → Disaggregated).
model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs Adds EPD encode planning + bootstrap/KV rendezvous injection for Chat API; emits an ExecutionPlan.
model_gateway/src/routers/grpc/proto_wrapper.rs Adds encode bootstrap info structs + TokenSpeed KV bootstrap setters; refines multimodal pixel stripping for TokenSpeed.
model_gateway/src/routers/grpc/pipeline.rs Builds new EPD pipelines and unifies execution stage to consume ExecutionPlan.
model_gateway/src/routers/grpc/pd_router.rs Adds EPD router constructor (new_epd) sharing PD routing logic but using EPD pipelines + router_type label.
model_gateway/src/routers/grpc/multimodal.rs Adds “after encode” multimodal assembly and TokenSpeed wire dtype / SHM handling updates for EPD.
model_gateway/src/routers/grpc/mod.rs Exposes new epd_encode module.
model_gateway/src/routers/grpc/harmony/streaming.rs Renames “Dual” streaming terminology to “PrefillDecode” in Harmony path.
model_gateway/src/routers/grpc/harmony/stages/request_building.rs Emits ExecutionPlan for Harmony requests and threads ExecutionPlanKind.
model_gateway/src/routers/grpc/epd_encode.rs New: TokenSpeed-specific encode dispatch planning and encode RPC dispatch implementation.
model_gateway/src/routers/grpc/context.rs Introduces ExecutionPlan/ExecutionPlanKind; renames PD selection variants to Disaggregated and adds encode assignments.
model_gateway/src/routers/grpc/common/stages/worker_selection.rs Adds EPD worker selection (encode-per-item + prefill/decode) and routing-key hashing for encode assignment.
model_gateway/src/routers/grpc/common/stages/request_execution.rs Refactors request execution to execute ExecutionPlan, including spawning encode dispatch jobs for EPD.
model_gateway/src/routers/grpc/common/stages/mod.rs Removes ExecutionMode export; request execution stage now plan-driven.
model_gateway/src/routers/grpc/common/stages/helpers.rs Adds EPD encode planning helper and TokenSpeed KV rendezvous injector.
model_gateway/src/routers/grpc/common/stages/dispatch_metadata.rs Updates dispatch metadata extraction to use ExecutionPlan.
model_gateway/src/routers/grpc/common/stages/client_acquisition.rs Renames client selection variant to Disaggregated.
model_gateway/src/routers/grpc/common/response_collection.rs Renames Dual → PrefillDecode in execution result handling/docs.
model_gateway/src/routers/factory.rs Adds GRPC_EPD router id and factory wiring for EPD router/policies.
model_gateway/src/policies/registry.rs Adds encode policy storage + default/fallback selection for EPD.
model_gateway/src/observability/metrics.rs Adds WORKER_ENCODE label and minor doc tweak.
model_gateway/src/main.rs Adds CLI flags for --epd-disaggregation, --encode, encode selector/policy, and routes into router/discovery config.
model_gateway/src/health.rs Adds readiness logic for EPD: require encode+prefill+decode.
model_gateway/src/config/validation.rs Validates EPD urls/ports/selectors and encode policy restrictions; adds tests.
model_gateway/src/config/types.rs Adds RoutingMode::EncodePrefillDecode, encode selector in discovery config, and routing-mode helpers.
grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py Adds TokenSpeed-side handling for encode bootstrap + KV bootstrap in Generate requests; adjusts multimodal parsing for EPD.
grpc_servicer/smg_grpc_servicer/tokenspeed/server.py Adds TokenSpeed encode role server composition (encoder servicer + scheduler servicer for discovery) and changes keepalive options/warmup behavior.
grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py New: TokenSpeed EPD encode gRPC servicer implementation.
crates/protocols/src/worker.rs Adds WorkerType::Encode to shared protocol enum parsing/printing.
crates/grpc_client/src/tokenspeed_scheduler.rs Initializes new TokenSpeed EPD bootstrap fields in scheduler request proto conversion.
crates/grpc_client/src/tokenspeed_encoder.rs New: cached/pool-based gRPC client for TokenSpeed encode service.
crates/grpc_client/src/lib.rs Exports the new TokenSpeed encode client/proto module.
crates/grpc_client/python/smg_grpc_proto/init.py Exposes generated Python stubs for TokenSpeed encoder service.
crates/grpc_client/proto/tokenspeed_scheduler.proto Extends scheduler GenerateRequest with encode + KV bootstrap info messages.
crates/grpc_client/proto/tokenspeed_encoder.proto New: TokenSpeed EPD encode service proto definition.
crates/grpc_client/build.rs Adds proto build pass for TokenSpeed encoder with scheduler extern_path mapping.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

};

/// gRPC PD (Prefill-Decode) router implementation for SGLang
/// gRPC disaggregated prefill/decode router router. Serves both PD (prefill-decode) and
Comment on lines 358 to +364
match worker.worker_type() {
WorkerType::Prefill => prefill_count += 1,
WorkerType::Decode => decode_count += 1,
WorkerType::Regular => regular_count += 1,
// EPD encode workers are a distinct pool: counted in `total` and
// listed, but not in the P/D/Regular sub-counts.
WorkerType::Encode => {}
Comment on lines +71 to +76
# Long EPD requests can spend minutes without sending DATA frames while
# the Rust client still sends HTTP/2 keepalive pings.
("grpc.http2.min_recv_ping_interval_without_data_ms", 10000),
("grpc.http2.max_pings_without_data", 0),
("grpc.http2.max_ping_strikes", 0),
("grpc.keepalive_permit_without_calls", True),
Comment on lines +99 to +109
if server_args.disaggregation_mode == "encode":
# EPD encode worker: serve the vision-only encode loop via the encoder
# service. ALSO mount the scheduler service so the gateway's generic
# worker discovery (HealthCheck + GetModelInfo + GetServerInfo, all over
# the TokenSpeedScheduler stub) can reach this worker and register it.
# Note: mounting the scheduler service means the unconditional startup
# warmup (stub.Generate) would route to the scheduler Generate path and
# drive the LM, which can SIGUSR1-kill the encode TP group -- encode
# workers must run with TOKENSPEED_SKIP_GRPC_WARMUP=1 (or warmup must be
# guarded to skip in encode mode). HealthCheck short-circuits to a
# shallow probe for disaggregation roles, so it won't drive the LM.

@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: 59d5010a01

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py Outdated
Comment thread model_gateway/src/routers/grpc/common/stages/worker_selection.rs Outdated
@chenht2022
chenht2022 requested a review from gongwei-130 as a code owner June 28, 2026 14:46
@github-actions github-actions Bot added the python-bindings Python bindings changes label Jun 28, 2026
@mergify

mergify Bot commented Jun 28, 2026

Copy link
Copy Markdown
Contributor

Hi @chenht2022, the DCO sign-off check has failed. All commits must include a Signed-off-by line.

To fix existing commits:

# Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-lease

To sign off future commits automatically:

  • Use git commit -s every time, or
  • VSCode: enable Git: Always Sign Off in Settings
  • PyCharm: enable Sign-off commit in the Commit tool window

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

Caution

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

⚠️ Outside diff range comments (2)
model_gateway/src/main.rs (1)

1545-1562: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚖️ Poor tradeoff

Argv filtering duplicates the optional-port consumption rule in parse_url_port_args.

The token-consumption logic here (consume <flag> <url> plus an optional trailing port/none) re-implements the exact rule in parse_url_port_args (Lines 35-49). They are consistent today, but any future change to one (e.g. accepting null or a different port format) silently desyncs argv stripping from URL parsing, causing flags to leak into clap or ports to be mis-bound. Consider extracting a single shared predicate for "does args[i+2] count as a consumed bootstrap-port token".

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/main.rs` around lines 1545 - 1562, The argv filtering in
main currently duplicates the optional bootstrap-port consumption rules already
implemented in parse_url_port_args, so keep both paths in sync by extracting a
shared predicate/helper for the trailing token check. Update the filtering loop
around filtered_args/raw_args to reuse that shared logic for consuming the
optional port/none token after --prefill or --encode, matching the behavior used
in parse_url_port_args so future format changes only need one fix.
model_gateway/src/routers/grpc/regular/streaming.rs (1)

654-681: 🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Drain or supervise prefill streams even when logprobs are disabled.

Both chat and generate PD helpers skip polling prefill_stream on the common no-logprobs path, so prefill RPC errors are ignored and the stream can remain unread until decode finishes. Drain it unconditionally, or spawn a supervised drain task if decode overlap must be preserved.

Proposed localized fix
-        if original_request.logprobs {
-            while let Some(response) = prefill_stream.next().await {
-                let gen_response =
-                    response.map_err(|e| format!("Prefill stream error: {}", e.message()))?;
-                match gen_response.into_response() {
-                    ProtoResponseVariant::Complete(_complete) => {
-                        // Input logprobs collected but not yet used in streaming
-                        // (OpenAI spec doesn't require prompt logprobs in streaming responses)
-                        break;
-                    }
-                    _ => continue,
-                }
-            }
-        }
+        while let Some(response) = prefill_stream.next().await {
+            let gen_response =
+                response.map_err(|e| format!("Prefill stream error: {}", e.message()))?;
+            if matches!(gen_response.into_response(), ProtoResponseVariant::Complete(_)) {
+                break;
+            }
+        }
-        let input_token_logprobs = if ctx.return_logprob {
-            let mut input_logprobs = None;
-            while let Some(response) = prefill_stream.next().await {
-                let gen_response =
-                    response.map_err(|e| format!("Prefill stream error: {}", e.message()))?;
-                match gen_response.into_response() {
-                    ProtoResponseVariant::Complete(complete) => {
-                        // Extract input_logprobs from prefill Complete message (convert proto to SGLang format)
-                        input_logprobs = complete
-                            .input_logprobs()
-                            .as_ref()
-                            .map(utils::convert_generate_input_logprobs);
-                        break;
-                    }
-                    _ => continue,
-                }
-            }
-            input_logprobs
-        } else {
-            None
-        };
+        let mut input_token_logprobs = None;
+        while let Some(response) = prefill_stream.next().await {
+            let gen_response =
+                response.map_err(|e| format!("Prefill stream error: {}", e.message()))?;
+            if let ProtoResponseVariant::Complete(complete) = gen_response.into_response() {
+                if ctx.return_logprob {
+                    input_token_logprobs = complete
+                        .input_logprobs()
+                        .as_ref()
+                        .map(utils::convert_generate_input_logprobs);
+                }
+                break;
+            }
+        }

Also applies to: 907-937

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/regular/streaming.rs` around lines 654 - 681,
The prefill stream in process_prefill_decode_streaming_chunks is only consumed
when original_request.logprobs is true, so the common no-logprobs path can
ignore prefill RPC errors and leave the stream unread. Update this PD helper to
always drain or supervise prefill_stream before/while decode proceeds, and apply
the same pattern to the related chat/generate PD helper mentioned in the review
so the prefill side is always polled and any errors are surfaced.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@model_gateway/src/config/validation.rs`:
- Around line 841-872: The EPD validation in validation.rs only checks explicit
prefill_policy and decode_policy overrides, so it can miss the effective policy
inherited from config.policy. Update the RoutingMode::EncodePrefillDecode
validation to resolve the effective prefill and decode policies first (using the
fallback from config.policy when the stage-specific policy is None), then apply
the PowerOfTwo worker-count checks and the Bucket ban against those effective
policies. Keep the fix localized around the existing
RoutingMode::EncodePrefillDecode branch and its prefill_policy/decode_policy
handling.

In `@model_gateway/src/routers/grpc/common/stages/request_execution.rs`:
- Around line 279-283: The EPD prefill path is still carrying raw multimodal
pixel tensors because `prefill_request` is cloned before
`decode_request.clear_mm_pixel_values()` and the PD helper cannot distinguish
EPD from non-EPD. Update the request flow in `execute_disaggregated_dispatch`
and the PD helper it calls to accept an EPD flag, then clear
`prefill_request.clear_mm_pixel_values()` only when that flag is set so encode
workers remain the sole owners of the vision-tower inputs.
- Around line 286-323: The encode dispatch path in
RequestExecutionStage::spawn_encode_dispatch currently only logs failures and
never records them against the selected worker, so a bad encode worker can keep
being chosen. Update the PreparedEncodeJob::dispatch flow (or the dispatch
caller in spawn_encode_dispatch) to call the same worker outcome/circuit-breaker
recording used by the prefill/decode selection paths whenever dispatch returns
Err or the spawned task panics, keeping the failure tied to the worker selection
path.

In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs`:
- Around line 494-509: `SelectWorkerInfo` now requires `leg`, so update the
`SelectWorkerInfo` construction used for prefill/decode and the one passed
through `assign_encode_workers` to include the correct leg value. Also route
these worker selections through the registry path instead of calling
`prefill_policy.select_worker` and `decode_policy.select_worker` directly, so
the `X-SMG-Routing-Key` sticky override is honored consistently in
`worker_selection.rs`.

In `@model_gateway/src/routers/grpc/context.rs`:
- Around line 326-328: Attach load guards to encode dispatch in the
disaggregated context so encode traffic is accounted for by load-aware routing
and overload protection. Update the `LoadGuards::Disaggregated` flow and the
`spawn_encode_dispatch()` path in `context.rs` to acquire and pass a
`WorkerLoadGuard` for each encode RPC, similar to how prefill/decode guards are
handled, while keeping the per-item EPD dispatch behavior intact.

In `@model_gateway/src/routers/grpc/multimodal.rs`:
- Around line 958-964: The SHM destination selection in
assemble_multimodal_data_after_encode() is using the wrong worker set for the
prefill leg, so local encode plus remote prefill can end up with unreadable
shared-memory handles. Update the assembly flow in assemble_multimodal_data_impl
and the related call sites to choose the SHM destination based on the current
assembly leg: after-encode should validate/target prefill workers, while
pre-encode should continue using encode workers. Use the existing
worker_shares_dev_shm() and encode_assignments logic, but make it branch on the
leg-specific worker selection instead of always checking encode workers.

In `@model_gateway/src/service_discovery.rs`:
- Line 318: The bootstrap-port check in service discovery is moving pod_type
because matches!(pod_type, ...) consumes the Option, which prevents later use
when constructing PodInfo. Update the bootstrap_port logic in the relevant
service_discovery code path to borrow pod_type in the matches! check (or
otherwise avoid taking ownership) so pod_type remains available for the
subsequent PodInfo build.

---

Outside diff comments:
In `@model_gateway/src/main.rs`:
- Around line 1545-1562: The argv filtering in main currently duplicates the
optional bootstrap-port consumption rules already implemented in
parse_url_port_args, so keep both paths in sync by extracting a shared
predicate/helper for the trailing token check. Update the filtering loop around
filtered_args/raw_args to reuse that shared logic for consuming the optional
port/none token after --prefill or --encode, matching the behavior used in
parse_url_port_args so future format changes only need one fix.

In `@model_gateway/src/routers/grpc/regular/streaming.rs`:
- Around line 654-681: The prefill stream in
process_prefill_decode_streaming_chunks is only consumed when
original_request.logprobs is true, so the common no-logprobs path can ignore
prefill RPC errors and leave the stream unread. Update this PD helper to always
drain or supervise prefill_stream before/while decode proceeds, and apply the
same pattern to the related chat/generate PD helper mentioned in the review so
the prefill side is always polled and any errors are surfaced.
🪄 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: de4d0677-6c54-4ae4-bf86-ee62dfd3c9da

📥 Commits

Reviewing files that changed from the base of the PR and between 9ec311c and 59d5010.

📒 Files selected for processing (49)
  • crates/grpc_client/build.rs
  • crates/grpc_client/proto/tokenspeed_encoder.proto
  • crates/grpc_client/proto/tokenspeed_scheduler.proto
  • crates/grpc_client/python/smg_grpc_proto/__init__.py
  • crates/grpc_client/src/lib.rs
  • crates/grpc_client/src/tokenspeed_encoder.rs
  • crates/grpc_client/src/tokenspeed_scheduler.rs
  • crates/protocols/src/worker.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/server.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
  • model_gateway/src/config/types.rs
  • model_gateway/src/config/validation.rs
  • model_gateway/src/health.rs
  • model_gateway/src/main.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/policies/registry.rs
  • model_gateway/src/routers/factory.rs
  • model_gateway/src/routers/grpc/common/response_collection.rs
  • model_gateway/src/routers/grpc/common/stages/client_acquisition.rs
  • model_gateway/src/routers/grpc/common/stages/dispatch_metadata.rs
  • model_gateway/src/routers/grpc/common/stages/helpers.rs
  • model_gateway/src/routers/grpc/common/stages/mod.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs
  • model_gateway/src/routers/grpc/context.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/harmony/stages/request_building.rs
  • model_gateway/src/routers/grpc/harmony/streaming.rs
  • model_gateway/src/routers/grpc/mod.rs
  • model_gateway/src/routers/grpc/multimodal.rs
  • model_gateway/src/routers/grpc/pd_router.rs
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/proto_wrapper.rs
  • model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/classify/response_processing.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/embedding/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/generate/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/request_building.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs
  • model_gateway/src/routers/router_manager.rs
  • model_gateway/src/service_discovery.rs
  • model_gateway/src/worker/manager.rs
  • model_gateway/src/worker/registry.rs
  • model_gateway/src/worker/service.rs
  • model_gateway/src/worker/worker.rs
  • model_gateway/src/workflow/job_queue.rs

Comment thread model_gateway/src/config/validation.rs
Comment thread model_gateway/src/routers/grpc/common/stages/request_execution.rs
Comment thread model_gateway/src/routers/grpc/common/stages/request_execution.rs Outdated
Comment thread model_gateway/src/routers/grpc/common/stages/worker_selection.rs Outdated
Comment thread model_gateway/src/routers/grpc/context.rs
Comment thread model_gateway/src/routers/grpc/multimodal.rs Outdated
Comment thread model_gateway/src/service_discovery.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

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/common/stages/worker_selection.rs (1)

438-473: 🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Select a runtime that has all required EPD legs.

Line 441 anchors EPD to the first prefill worker’s runtime. If that runtime lacks encode or decode workers, this returns model_not_found even when another runtime has encode+prefill+decode available. Choose a runtime from the intersection of available prefill/decode/(encode when needed) pools before filtering.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs` around
lines 438 - 473, The runtime selection in worker selection is anchored to the
first prefill worker, which can reject valid EPD combinations when that runtime
lacks one of the required legs. Update the selection logic in the EPD path to
choose a runtime present in the intersection of the available prefill, decode,
and (when needed) encode pools before filtering. Keep the mixed-runtime warning
and subsequent filtering in the same flow, but base them on the chosen runtime
rather than all_prefill.first().
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@bindings/python/src/lib.rs`:
- Line 610: The Python bindings currently hardcode encode_selector to an empty
map, so encode worker discovery cannot be configured. Update the relevant PyO3
constructor/build path in lib.rs to add an encode_selector field/argument
alongside prefill_selector and decode_selector, preserve positional
compatibility, and pass/clone it through both config builders instead of using
HashMap::new(). Make sure the same encode_selector value is wired through the
TokenSpeed EPD discovery setup so Python callers can control encode-role worker
selection.

---

Outside diff comments:
In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs`:
- Around line 438-473: The runtime selection in worker selection is anchored to
the first prefill worker, which can reject valid EPD combinations when that
runtime lacks one of the required legs. Update the selection logic in the EPD
path to choose a runtime present in the intersection of the available prefill,
decode, and (when needed) encode pools before filtering. Keep the mixed-runtime
warning and subsequent filtering in the same flow, but base them on the chosen
runtime rather than all_prefill.first().
🪄 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: 66098e17-9a15-4696-991e-a19b35bc8489

📥 Commits

Reviewing files that changed from the base of the PR and between 59d5010 and 9fe7622.

📒 Files selected for processing (2)
  • bindings/python/src/lib.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs

Comment thread bindings/python/src/lib.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: 9fe7622430

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +326 to 330
/// Disaggregated guards cover the prefill+decode pair. EPD encode workers are
/// assigned per item; their fire-and-supervise RPCs do not hold load guards.
Disaggregated {
_prefill: WorkerLoadGuard,
_decode: WorkerLoadGuard,

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 Track encode worker load for least-load routing

When encode_policy is least_load and encode workers have no fresh GetLoads snapshot (cold start, or encode workers that do not report scheduler loads), LeastLoadPolicy falls back to Worker::load() to avoid sending all traffic to the first worker. This new EPD guard set only increments prefill/decode and explicitly leaves encode jobs unguarded, so encode worker load remains 0 while per-item encode RPCs are in flight; concurrent image requests therefore all select the same encode worker until external load reports arrive. Hold a guard for assigned encode workers, or reject least_load for encode, so the advertised policy can balance the encode fan-out.

Useful? React with 👍 / 👎.

@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: 54594f006a

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment thread model_gateway/src/routers/grpc/multimodal.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

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/config/validation.rs (1)

808-885: 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Reject EPD decode Bucket even with service discovery enabled.

The worker-count checks depend on static URLs, but the decode Bucket ban is policy compatibility and should not be skipped when discovery is enabled. As written, EPD service-discovery deployments can pass validation with decode_policy = Bucket or a fallback main Bucket.

Proposed fix
-        if !has_service_discovery {
+        if let RoutingMode::EncodePrefillDecode {
+            decode_policy, ..
+        } = &config.mode
+        {
+            let effective_decode_policy = decode_policy.as_ref().unwrap_or(&config.policy);
+            if matches!(effective_decode_policy, PolicyConfig::Bucket { .. }) {
+                return Err(ConfigError::IncompatibleConfig {
+                    reason: "Decode policy should not be allowed to be bucket".to_string(),
+                });
+            }
+        }
+
+        if !has_service_discovery {
             if let PolicyConfig::PowerOfTwo { .. } = &config.policy {

Then remove the inner duplicate Bucket check at lines 880-884.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/config/validation.rs` around lines 808 - 885, The
validation in the RoutingMode::EncodePrefillDecode branch is skipping the decode
Bucket compatibility rule when service discovery is enabled, so decode_policy or
the fallback config.policy can still be Bucket. Move the decode Bucket rejection
out of the !has_service_discovery guard so it always applies, and keep the
existing effective_decode_policy logic in validation.rs to check both the
explicit decode_policy and the fallback main policy. Remove the duplicate inner
Bucket check from EncodePrefillDecode after adding the unconditional validation.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py`:
- Around line 199-202: Handle failures in the encoder servicer flow explicitly
so RPCs return the intended gRPC status instead of UNKNOWN: wrap the `payload =
await asyncio.to_thread(self._parse_and_pickle, request, bootstrap_room)` path
and the scheduler/ingest call in the relevant servicer method in a try/except,
and in the exception handler log the exception and call
`context.abort(grpc.StatusCode.INTERNAL, str(e))` without altering the exception
message. Keep the fix localized to the code path around `_parse_and_pickle`,
`send_to_scheduler.send`, and `_ingest` so malformed input and backend errors
are surfaced consistently.

---

Outside diff comments:
In `@model_gateway/src/config/validation.rs`:
- Around line 808-885: The validation in the RoutingMode::EncodePrefillDecode
branch is skipping the decode Bucket compatibility rule when service discovery
is enabled, so decode_policy or the fallback config.policy can still be Bucket.
Move the decode Bucket rejection out of the !has_service_discovery guard so it
always applies, and keep the existing effective_decode_policy logic in
validation.rs to check both the explicit decode_policy and the fallback main
policy. Remove the duplicate inner Bucket check from EncodePrefillDecode after
adding the unconditional validation.
🪄 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: 92e15c4d-224c-4a84-be6b-4c44f6715628

📥 Commits

Reviewing files that changed from the base of the PR and between 9fe7622 and 6592f9d.

📒 Files selected for processing (8)
  • bindings/python/src/lib.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/server.py
  • model_gateway/src/config/validation.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/multimodal.rs

Comment thread grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py 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: 35d2caeacd

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

// We need at least one prefill and one decode worker to handle requests
// in PD disaggregation mode.
let grpc_pd = if grpc_prefill > 0 && grpc_decode > 0 {
let grpc_pd = if !grpc_epd_ready && grpc_prefill > 0 && grpc_decode > 0 {

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 PD router eligible when EPD workers exist

In IGW/service-discovery deployments that have gRPC PD workers plus an encode worker for the same model (for example during an EPD rollout), this condition zeros the GRPC_PD weight as soon as grpc_epd_ready is true, forcing PD-capable traffic through the EPD router. The EPD chat/generate pipeline is built with inject_pd_metadata=false, so if the selected prefill/decode runtime is SGLang for a text-only request, the request reaches the PD workers without the required disaggregated bootstrap metadata and can fail or hang; keep the PD weight when PD workers exist, or make the EPD-ready path only count TokenSpeed EPD worker sets.

Useful? React with 👍 / 👎.

(
router_ids::GRPC_EPD,
"gRPC EPD",
Self::create_grpc_epd_router(None, None, None, policy, ctx),

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 Honor EPD per-role policies in IGW

When Kubernetes service discovery is enabled, main auto-enables IGW, so EPD deployments use this factory path instead of the single-router path. Passing None for encode/prefill/decode discards the RoutingMode::EncodePrefillDecode policies parsed from CLI/config (for example --encode-policy random or per-leg prefill/decode policies), causing the EPD router to silently use the default encode policy and the main policy for P/D; thread the mode's per-role policies into IGW router creation.

Useful? React with 👍 / 👎.

@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: 24cac7b47b

ℹ️ About Codex in GitHub

Your team has set up Codex to 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 👍.

Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".

Comment on lines +172 to +175
RoutingMode::EncodePrefillDecode { .. } => {
let has_encode = healthy_workers
.iter()
.any(|w| matches!(w.worker_type(), WorkerType::Encode));

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 Gate IGW readiness on complete EPD worker sets

When service discovery is enabled, main auto-enables IGW, so this EPD role check is skipped and workers_ready becomes just !healthy_workers.is_empty(). In an EPD rollout where the encode pod becomes healthy before prefill/decode, /readiness can report ready even though no router can serve the model yet, causing kube/load balancers to send traffic that will fail until all E/P/D roles are present; apply the same complete-role gating to IGW when the only healthy workers are disaggregated auxiliaries.

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.

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/epd_encode.rs (1)

169-182: 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Prepare items before requiring encode assignments.

The empty-plan fast path at Lines 177-182 cannot run when encode_assignments() is None or empty, because Lines 169-174 error first. Reorder this so zero prepared encode items returns an empty plan without requiring encode worker assignments.

Proposed fix
-    let workers = workers.ok_or_else(|| anyhow!("Worker selection stage not completed"))?;
-    let encode_assignments = workers
-        .encode_assignments()
-        .filter(|assignments| !assignments.is_empty())
-        .ok_or_else(|| anyhow!("Encode planning requires EPD worker selection"))?
-        .to_vec();
-
-    let items = prepare_items(precomputed, clients, Some(workers))?;
+    let workers = workers.ok_or_else(|| anyhow!("Worker selection stage not completed"))?;
+    let items = prepare_items(precomputed, clients, Some(workers))?;
     if items.is_empty() {
         return Ok(EncodePlan {
             bootstrap_info: Vec::new(),
             dispatch: EncodeDispatchPlan::new(Vec::new()),
         });
     }
+
+    let encode_assignments = workers
+        .encode_assignments()
+        .filter(|assignments| !assignments.is_empty())
+        .ok_or_else(|| anyhow!("Encode planning requires EPD worker selection"))?
+        .to_vec();
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/epd_encode.rs` around lines 169 - 182, The
empty-plan fast path in epd_encode::encode_plan is blocked by the early
encode_assignments check, so zero prepared items still fail before returning an
empty EncodePlan. Reorder the logic so prepare_items runs first, then
immediately return the empty plan when items.is_empty() without calling
workers.encode_assignments(); only require and materialize encode_assignments
when there are items to dispatch.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Outside diff comments:
In `@model_gateway/src/routers/grpc/epd_encode.rs`:
- Around line 169-182: The empty-plan fast path in epd_encode::encode_plan is
blocked by the early encode_assignments check, so zero prepared items still fail
before returning an empty EncodePlan. Reorder the logic so prepare_items runs
first, then immediately return the empty plan when items.is_empty() without
calling workers.encode_assignments(); only require and materialize
encode_assignments when there are items to dispatch.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 65cc6661-895d-4db2-a4d6-b2a7fbe920e4

📥 Commits

Reviewing files that changed from the base of the PR and between 1497856 and c4e8500.

📒 Files selected for processing (11)
  • bindings/python/src/lib.rs
  • crates/grpc_client/src/tokenspeed_encoder.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • model_gateway/src/config/validation.rs
  • model_gateway/src/main.rs
  • model_gateway/src/routers/grpc/common/stages/helpers.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs
  • model_gateway/src/service_discovery.rs
💤 Files with no reviewable changes (1)
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs

@chenht2022

Copy link
Copy Markdown
Member Author

@coderabbitai resume

@coderabbitai

coderabbitai Bot commented Jul 6, 2026

Copy link
Copy Markdown
✅ Action performed

Reviews resumed.

@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/config/validation.rs (1)

243-246: 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

Validate the effective encode policy, not only explicit overrides.

Line 243 rejects unsupported encode policies only when encode_policy is set. If EPD encode selection inherits config.policy when this is None, top-level CacheAware/LeastLoad/Bucket can still reach the encode stage.

Proposed fix
-        if let RoutingMode::EncodePrefillDecode { decode_policy, .. } = &config.mode {
+        if let RoutingMode::EncodePrefillDecode {
+            encode_policy,
+            decode_policy,
+            ..
+        } = &config.mode
+        {
+            let effective_encode_policy = encode_policy.as_ref().unwrap_or(&config.policy);
+            Self::validate_encode_policy(effective_encode_policy)?;
             let effective_decode_policy = decode_policy.as_ref().unwrap_or(&config.policy);
             if matches!(effective_decode_policy, PolicyConfig::Bucket { .. }) {
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/config/validation.rs` around lines 243 - 246, The
encode-policy validation in the config validation flow only checks explicit
overrides, so inherited values from config.policy can slip through. Update the
validation logic in the encode-policy path of validation.rs to validate the
effective encode policy that will actually be used by EPD, not just the optional
encode_policy override. Use the existing validate_policy and
validate_encode_policy helpers from the same validation routine, but resolve the
fallback policy first so top-level CacheAware, LeastLoad, and Bucket values are
rejected before reaching encoding.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Outside diff comments:
In `@model_gateway/src/config/validation.rs`:
- Around line 243-246: The encode-policy validation in the config validation
flow only checks explicit overrides, so inherited values from config.policy can
slip through. Update the validation logic in the encode-policy path of
validation.rs to validate the effective encode policy that will actually be used
by EPD, not just the optional encode_policy override. Use the existing
validate_policy and validate_encode_policy helpers from the same validation
routine, but resolve the fallback policy first so top-level CacheAware,
LeastLoad, and Bucket values are rejected before reaching encoding.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: ASSERTIVE

Plan: Pro

Run ID: 70669df3-eeb8-495c-8f86-228a9ffcf857

📥 Commits

Reviewing files that changed from the base of the PR and between 1497856 and 5782193.

📒 Files selected for processing (11)
  • bindings/python/src/lib.rs
  • crates/grpc_client/src/tokenspeed_encoder.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • model_gateway/src/config/validation.rs
  • model_gateway/src/main.rs
  • model_gateway/src/routers/grpc/common/stages/helpers.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs
  • model_gateway/src/service_discovery.rs
💤 Files with no reviewable changes (1)
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs

@chenht2022

Copy link
Copy Markdown
Member Author

@coderabbitai I think this is a false positive. EPD encode does not inherit config.policy when encode_policy is unset; it defaults to ConsistentHashing in create_grpc_epd_router, and RoutingMode::get_encode_policy documents the same behavior. config.policy is only the fallback for prefill/decode. Using unwrap_or(&config.policy) here would incorrectly reject valid EPD configs where the main policy is cache_aware or least_load and encode intentionally uses its default.

@chenht2022

Copy link
Copy Markdown
Member Author

@coderabbitai approve

@coderabbitai

coderabbitai Bot commented Jul 6, 2026

Copy link
Copy Markdown
✅ Action performed

Comments resolved and changes approved.

@chenht2022

Copy link
Copy Markdown
Member Author

@coderabbitai review

@coderabbitai

coderabbitai Bot commented Jul 6, 2026 •

Copy link
Copy Markdown
✅ Action performed

Review finished.

Note: CodeRabbit is an incremental review system and does not re-review already reviewed commits. This command is applicable only when automatic reviews are paused.

@slin1237
slin1237 force-pushed the epd-routing-clean branch from 5782193 to 83db4b2 Compare July 7, 2026 12:34

@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: 83db4b2bb4

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

RoutingMode::Anthropic {
worker_urls: self.worker_urls.clone(),
}
} else if self.epd_disaggregation {

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 EPD mode when Python IGW is enabled

When Python callers enable IGW for service discovery, this EPD arm is never reached because the earlier enable_igw arm forces RoutingMode::Regular. The subsequent router_config.validate() therefore applies the Regular discovery rule, which requires selector, instead of the EPD rule using encode_selector/prefill_selector/decode_selector, so an EPD Kubernetes deployment built through the Python API cannot start in IGW even though the CLI keeps the EPD routing mode under IGW. Keep EncodePrefillDecode here and let IGW only affect router initialization.

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: 5

Caution

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

⚠️ Outside diff range comments (4)
model_gateway/src/service_discovery.rs (1)

403-441: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Minor duplication in selector-to-string formatting.

The encode/prefill/decode selector-stringification block repeats the same .iter().map(|(k,v)| format!("{k}={v}")).collect::<Vec<_>>().join(",") pattern already used in list_label_selector/router_label_selector. A small fn format_selector(m: &HashMap<String,String>) -> String helper would remove the repetition, but this is log-only formatting with no behavioral impact.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/service_discovery.rs` around lines 403 - 441, The selector
string formatting in service_discovery is duplicated across the disaggregated
and non-disaggregated logging paths. Extract the repeated
`.iter().map(...).collect::<Vec<_>>().join(",")` logic into a shared helper like
`format_selector`, and reuse it for `encode_selector`, `prefill_selector`,
`decode_selector`, and `selector` so the logging in `service_discovery` stays
consistent and easier to maintain.
model_gateway/src/worker/service.rs (1)

340-378: 🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

Add encode_count to ListWorkersResult. /workers leaves WorkerType::Encode out of its aggregate counters, while WorkerRegistry::stats() already reports encode_workers. Adding the field would keep the listing response consistent with the registry stats API.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/worker/service.rs` around lines 340 - 378, Add encode_count
to the worker listing response so ListWorkersResult matches
WorkerRegistry::stats(). Update the list_workers method in service.rs to track
WorkerType::Encode alongside prefill_count, decode_count, and regular_count, and
include that value in the returned ListWorkersResult. Make sure the /workers
output still counts Encode workers in total while keeping the existing per-type
counters unchanged for the other worker types.
model_gateway/src/routers/grpc/pipeline.rs (1)

259-333: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

new_pd/new_epd duplicate identical processor/streaming_processor construction.

Both functions build ResponseProcessor and StreamingProcessor with identical arguments (same metrics_labels::BACKEND_PD); only the WorkerSelectionMode and request-building stage's ExecutionPlanKind/inject_pd_metadata differ. The same pattern repeats for new_messages_pd/new_messages_epd (Lines 443-544) and new_completion_pd/new_completion_epd (Lines 597-691). Extracting the processor/streaming_processor setup into a small helper would remove ~15 duplicated lines per pair (×3 pairs) introduced by this PR.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/pipeline.rs` around lines 259 - 333, The
PD/EPD pipeline constructors duplicate the same `ResponseProcessor` and
`StreamingProcessor` setup, so factor that shared initialization into a helper
and reuse it from `new_pd`, `new_epd`, `new_messages_pd`, `new_messages_epd`,
`new_completion_pd`, and `new_completion_epd`. Keep the existing differences in
each constructor limited to `WorkerSelectionMode`, `ExecutionPlanKind`, and the
`inject_pd_metadata` flag in `ChatGenerateRequestBuildingStage::new`, while
moving the identical processor creation around `metrics_labels::BACKEND_PD` into
one shared function.
model_gateway/src/routers/grpc/harmony/streaming.rs (1)

583-586: 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

Guard prefill_stream.mark_completed() on decode success
model_gateway/src/routers/grpc/harmony/streaming.rs:585 should only mark the prefill stream completed when process_decode_stream() returns Ok(_). On a decode error, let the prefill stream drop so the backend aborts, matching the regular PD path and shared collection helper.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/harmony/streaming.rs` around lines 583 - 586,
In `process_decode_stream` handling, `prefill_stream.mark_completed()` is being
called unconditionally after the await, which incorrectly marks the prefill
stream complete even when decode fails. Update the logic in `streaming.rs` so
`mark_completed()` is only invoked when `process_decode_stream()` returns
`Ok(_)`, and leave the prefill stream to drop on errors to preserve the backend
abort behavior used by the regular PD path and shared collection helper.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@bindings/python/src/lib.rs`:
- Around line 584-604: Reject conflicting routing flags in the routing builder:
`epd_disaggregation` and `pd_disaggregation` must not both be enabled, because
the current `else if` chain in the `RoutingMode` selection silently drops the
`pd_disaggregation` configuration. Add an explicit validation check before
constructing `RoutingMode` in this logic (near `convert_policy` usage) and
return an error when both flags are true, so callers must choose either EPD or
PD and the active mode is unambiguous.

In `@model_gateway/src/config/validation.rs`:
- Around line 586-595: The EPD discovery validation in ConfigValidator::validate
currently enforces that encode_selector, prefill_selector, and decode_selector
are all non-empty; confirm this stricter requirement is intentional for
RoutingMode::EncodePrefillDecode and, if so, keep the branch as-is but add unit
coverage for it. Add a test alongside the existing validation tests that builds
a RouterConfig with DiscoveryConfig enabled and leaves one selector empty (for
example encode_selector) while the others are set, then assert validation fails;
mirror the style of test_validate_epd_mode_encode_policy_restrictions so the new
discovery-only rule is exercised directly.

In `@model_gateway/src/main.rs`:
- Around line 1023-1027: The CLI bucket policy currently hardcodes
bucket_adjust_interval_secs in PolicyConfig::Bucket, so users cannot tune it
from the command line. Add a matching CliArgs #[arg] field such as
--bucket-adjust-interval-secs, thread it through the bucket policy construction
in main.rs, and use that value instead of the fixed 5 to align CLI behavior with
the Python bindings’ bucket_adjust_interval_secs support.
- Around line 1219-1227: Add validation so `epd_disaggregation` and
`pd_disaggregation` cannot both be enabled at the same time in
`to_router_config` and the CLI parsing path in `main.rs`. Use clap
`conflicts_with` on the relevant flags or add an explicit guard before building
`RoutingMode::EncodePrefillDecode` / the PD routing mode, so one configuration
is rejected instead of silently taking precedence. Also update the startup
banner printing blocks around `pd_disaggregation` and `epd_disaggregation` to
only emit one set of node listings after the conflict is prevented.

In `@model_gateway/src/routers/grpc/pd_router.rs`:
- Around line 116-185: `GrpcPDRouter::new_epd` duplicates the shared constructor
setup already present in `GrpcPDRouter::new`. Extract the common initialization
logic for registry cloning, parser factory lookup, multimodal setup,
`SharedComponents` creation, and `retry_config` into a private helper, then have
both `new` and `new_epd` call it and only განსხვავate the pipeline builders and
`router_type` assignment.

---

Outside diff comments:
In `@model_gateway/src/routers/grpc/harmony/streaming.rs`:
- Around line 583-586: In `process_decode_stream` handling,
`prefill_stream.mark_completed()` is being called unconditionally after the
await, which incorrectly marks the prefill stream complete even when decode
fails. Update the logic in `streaming.rs` so `mark_completed()` is only invoked
when `process_decode_stream()` returns `Ok(_)`, and leave the prefill stream to
drop on errors to preserve the backend abort behavior used by the regular PD
path and shared collection helper.

In `@model_gateway/src/routers/grpc/pipeline.rs`:
- Around line 259-333: The PD/EPD pipeline constructors duplicate the same
`ResponseProcessor` and `StreamingProcessor` setup, so factor that shared
initialization into a helper and reuse it from `new_pd`, `new_epd`,
`new_messages_pd`, `new_messages_epd`, `new_completion_pd`, and
`new_completion_epd`. Keep the existing differences in each constructor limited
to `WorkerSelectionMode`, `ExecutionPlanKind`, and the `inject_pd_metadata` flag
in `ChatGenerateRequestBuildingStage::new`, while moving the identical processor
creation around `metrics_labels::BACKEND_PD` into one shared function.

In `@model_gateway/src/service_discovery.rs`:
- Around line 403-441: The selector string formatting in service_discovery is
duplicated across the disaggregated and non-disaggregated logging paths. Extract
the repeated `.iter().map(...).collect::<Vec<_>>().join(",")` logic into a
shared helper like `format_selector`, and reuse it for `encode_selector`,
`prefill_selector`, `decode_selector`, and `selector` so the logging in
`service_discovery` stays consistent and easier to maintain.

In `@model_gateway/src/worker/service.rs`:
- Around line 340-378: Add encode_count to the worker listing response so
ListWorkersResult matches WorkerRegistry::stats(). Update the list_workers
method in service.rs to track WorkerType::Encode alongside prefill_count,
decode_count, and regular_count, and include that value in the returned
ListWorkersResult. Make sure the /workers output still counts Encode workers in
total while keeping the existing per-type counters unchanged for the other
worker types.
🪄 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: d8f2fda6-38b2-4bfd-abfa-c7f49b84468b

📥 Commits

Reviewing files that changed from the base of the PR and between 5782193 and 83db4b2.

📒 Files selected for processing (50)
  • bindings/python/src/lib.rs
  • crates/grpc_client/build.rs
  • crates/grpc_client/proto/tokenspeed_encoder.proto
  • crates/grpc_client/proto/tokenspeed_scheduler.proto
  • crates/grpc_client/python/smg_grpc_proto/__init__.py
  • crates/grpc_client/src/lib.rs
  • crates/grpc_client/src/tokenspeed_encoder.rs
  • crates/grpc_client/src/tokenspeed_scheduler.rs
  • crates/protocols/src/worker.rs
  • grpc_servicer/smg_grpc_servicer/tokenspeed/encoder_servicer.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/server.py
  • grpc_servicer/smg_grpc_servicer/tokenspeed/servicer.py
  • model_gateway/src/config/types.rs
  • model_gateway/src/config/validation.rs
  • model_gateway/src/health.rs
  • model_gateway/src/main.rs
  • model_gateway/src/observability/metrics.rs
  • model_gateway/src/policies/registry.rs
  • model_gateway/src/routers/factory.rs
  • model_gateway/src/routers/grpc/common/response_collection.rs
  • model_gateway/src/routers/grpc/common/stages/client_acquisition.rs
  • model_gateway/src/routers/grpc/common/stages/dispatch_metadata.rs
  • model_gateway/src/routers/grpc/common/stages/helpers.rs
  • model_gateway/src/routers/grpc/common/stages/mod.rs
  • model_gateway/src/routers/grpc/common/stages/request_execution.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs
  • model_gateway/src/routers/grpc/context.rs
  • model_gateway/src/routers/grpc/epd_encode.rs
  • model_gateway/src/routers/grpc/harmony/stages/request_building.rs
  • model_gateway/src/routers/grpc/harmony/streaming.rs
  • model_gateway/src/routers/grpc/mod.rs
  • model_gateway/src/routers/grpc/multimodal.rs
  • model_gateway/src/routers/grpc/pd_router.rs
  • model_gateway/src/routers/grpc/pipeline.rs
  • model_gateway/src/routers/grpc/proto_wrapper.rs
  • model_gateway/src/routers/grpc/regular/stages/chat/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/classify/response_processing.rs
  • model_gateway/src/routers/grpc/regular/stages/completion/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/embedding/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/generate/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/messages/request_building.rs
  • model_gateway/src/routers/grpc/regular/stages/request_building.rs
  • model_gateway/src/routers/grpc/regular/streaming.rs
  • model_gateway/src/routers/router_manager.rs
  • model_gateway/src/service_discovery.rs
  • model_gateway/src/worker/manager.rs
  • model_gateway/src/worker/registry.rs
  • model_gateway/src/worker/service.rs
  • model_gateway/src/worker/worker.rs
  • model_gateway/src/workflow/job_queue.rs

Comment on lines +584 to +604
} else if self.epd_disaggregation {
RoutingMode::EncodePrefillDecode {
encode_urls: self.encode_urls.clone().unwrap_or_default(),
prefill_urls: self.prefill_urls.clone().unwrap_or_default(),
decode_urls: self.decode_urls.clone().unwrap_or_default(),
encode_policy: self
.encode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
prefill_policy: self
.prefill_policy
.as_ref()
.map(convert_policy)
.transpose()?,
decode_policy: self
.decode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win

No validation prevents epd_disaggregation and pd_disaggregation from both being true.

If a caller sets both flags, epd_disaggregation silently takes precedence (the else if chain) and pd_disaggregation's prefill/decode policy configuration is dropped without any error/warning. This is an easy misconfiguration to make (e.g., a caller migrating from PD to EPD and forgetting to unset the old flag) that would silently change routing behavior in production.

🐛 Suggested fix: reject the conflicting configuration explicitly
         } else if self.epd_disaggregation {
+            if self.pd_disaggregation {
+                return Err(config::ConfigError::IncompatibleConfig {
+                    reason: "epd_disaggregation and pd_disaggregation cannot both be enabled"
+                        .to_string(),
+                });
+            }
             RoutingMode::EncodePrefillDecode {
📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
} else if self.epd_disaggregation {
RoutingMode::EncodePrefillDecode {
encode_urls: self.encode_urls.clone().unwrap_or_default(),
prefill_urls: self.prefill_urls.clone().unwrap_or_default(),
decode_urls: self.decode_urls.clone().unwrap_or_default(),
encode_policy: self
.encode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
prefill_policy: self
.prefill_policy
.as_ref()
.map(convert_policy)
.transpose()?,
decode_policy: self
.decode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
}
} else if self.epd_disaggregation {
if self.pd_disaggregation {
return Err(config::ConfigError::IncompatibleConfig {
reason: "epd_disaggregation and pd_disaggregation cannot both be enabled"
.to_string(),
});
}
RoutingMode::EncodePrefillDecode {
encode_urls: self.encode_urls.clone().unwrap_or_default(),
prefill_urls: self.prefill_urls.clone().unwrap_or_default(),
decode_urls: self.decode_urls.clone().unwrap_or_default(),
encode_policy: self
.encode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
prefill_policy: self
.prefill_policy
.as_ref()
.map(convert_policy)
.transpose()?,
decode_policy: self
.decode_policy
.as_ref()
.map(convert_policy)
.transpose()?,
}
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@bindings/python/src/lib.rs` around lines 584 - 604, Reject conflicting
routing flags in the routing builder: `epd_disaggregation` and
`pd_disaggregation` must not both be enabled, because the current `else if`
chain in the `RoutingMode` selection silently drops the `pd_disaggregation`
configuration. Add an explicit validation check before constructing
`RoutingMode` in this logic (near `convert_policy` usage) and return an error
when both flags are true, so callers must choose either EPD or PD and the active
mode is unambiguous.

Comment thread model_gateway/src/config/validation.rs
Comment thread model_gateway/src/main.rs
Comment on lines +1023 to +1027
"bucket" => PolicyConfig::Bucket {
balance_abs_threshold: self.balance_abs_threshold,
balance_rel_threshold: self.balance_rel_threshold,
bucket_adjust_interval_secs: 5,
},

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick win

bucket_adjust_interval_secs is hardcoded to 5 for the CLI, with no way to configure it.

Unlike balance_abs_threshold/balance_rel_threshold (reused CLI flags), bucket_adjust_interval_secs has no corresponding #[arg] field in CliArgs, so CLI users can never tune it — while the Python bindings (bindings/python/src/lib.rs) already expose a real bucket_adjust_interval_secs field for this. Since this PR is the first to correctly wire up the "bucket" policy for the CLI (previously it silently fell back to RoundRobin), consider adding a matching --bucket-adjust-interval-secs flag for feature parity.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/main.rs` around lines 1023 - 1027, The CLI bucket policy
currently hardcodes bucket_adjust_interval_secs in PolicyConfig::Bucket, so
users cannot tune it from the command line. Add a matching CliArgs #[arg] field
such as --bucket-adjust-interval-secs, thread it through the bucket policy
construction in main.rs, and use that value instead of the fixed 5 to align CLI
behavior with the Python bindings’ bucket_adjust_interval_secs support.

Comment thread model_gateway/src/main.rs
Comment on lines +116 to +185
/// Create a new gRPC EPD (encode-prefill-decode) router.
///
/// Identical to [`Self::new`] except the chat/generate and messages
/// pipelines plan encode-worker jobs during request building, and all
/// pipelines dispatch through `ExecutionPlan::EncodePrefillDecode`.
/// Completion is text-only so it has no encode jobs but still uses the EPD
/// execution path.
pub fn new_epd(ctx: &Arc<AppContext>) -> Result<Self, String> {
let worker_registry = ctx.worker_registry.clone();
let policy_registry = ctx.policy_registry.clone();
let tokenizer_registry = ctx.tokenizer_registry.clone();

let reasoning_parser_factory = ctx
.reasoning_parser_factory
.as_ref()
.ok_or_else(|| "gRPC EPD router requires reasoning parser factory".to_string())?
.clone();
let tool_parser_factory = ctx
.tool_parser_factory
.as_ref()
.ok_or_else(|| "gRPC EPD router requires tool parser factory".to_string())?
.clone();

let multimodal = match MultimodalComponents::new(ctx.multimodal_config_registry.clone()) {
Ok(mc) => Some(Arc::new(mc)),
Err(e) => {
tracing::warn!("Multimodal components initialization failed (non-fatal): {e}");
None
}
};

let shared_components = Arc::new(SharedComponents {
tokenizer_registry: tokenizer_registry.clone(),
tool_parser_factory: tool_parser_factory.clone(),
reasoning_parser_factory: reasoning_parser_factory.clone(),
configured_tool_parser: ctx.configured_tool_parser.clone(),
multimodal,
});

let pipeline = RequestPipeline::new_epd(
worker_registry.clone(),
policy_registry.clone(),
tool_parser_factory.clone(),
reasoning_parser_factory.clone(),
ctx.configured_tool_parser.clone(),
ctx.configured_reasoning_parser.clone(),
);

let messages_pipeline = RequestPipeline::new_messages_epd(
worker_registry.clone(),
policy_registry.clone(),
tool_parser_factory.clone(),
reasoning_parser_factory.clone(),
ctx.configured_tool_parser.clone(),
ctx.configured_reasoning_parser.clone(),
);

let completion_pipeline =
RequestPipeline::new_completion_epd(worker_registry.clone(), policy_registry.clone());

Ok(GrpcPDRouter {
worker_registry,
pipeline,
messages_pipeline,
completion_pipeline,
shared_components,
retry_config: ctx.router_config.effective_retry_config(),
router_type: "grpc_epd",
})
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Extract shared setup logic between new and new_epd.

new_epd duplicates nearly the entire body of new (tokenizer/tool/reasoning factory resolution, multimodal component init, SharedComponents construction, retry_config init) — only the pipeline-builder calls and router_type value differ. Consider extracting a private helper that returns (worker_registry, policy_registry, shared_components, retry_config) and have both constructors call it before building their pipelines.

♻️ Sketch of extracted helper
+    fn build_common(ctx: &Arc<AppContext>) -> Result<(Arc<WorkerRegistry>, Arc<PolicyRegistry>, Arc<SharedComponents>, ToolParserFactory, ReasoningParserFactory), String> {
+        let worker_registry = ctx.worker_registry.clone();
+        let policy_registry = ctx.policy_registry.clone();
+        let tokenizer_registry = ctx.tokenizer_registry.clone();
+        let reasoning_parser_factory = ctx.reasoning_parser_factory.as_ref()
+            .ok_or_else(|| "gRPC router requires reasoning parser factory".to_string())?.clone();
+        let tool_parser_factory = ctx.tool_parser_factory.as_ref()
+            .ok_or_else(|| "gRPC router requires tool parser factory".to_string())?.clone();
+        let multimodal = match MultimodalComponents::new(ctx.multimodal_config_registry.clone()) {
+            Ok(mc) => Some(Arc::new(mc)),
+            Err(e) => { tracing::warn!("Multimodal components initialization failed (non-fatal): {e}"); None }
+        };
+        let shared_components = Arc::new(SharedComponents {
+            tokenizer_registry, tool_parser_factory: tool_parser_factory.clone(),
+            reasoning_parser_factory: reasoning_parser_factory.clone(),
+            configured_tool_parser: ctx.configured_tool_parser.clone(), multimodal,
+        });
+        Ok((worker_registry, policy_registry, shared_components, tool_parser_factory, reasoning_parser_factory))
+    }
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/pd_router.rs` around lines 116 - 185,
`GrpcPDRouter::new_epd` duplicates the shared constructor setup already present
in `GrpcPDRouter::new`. Extract the common initialization logic for registry
cloning, parser factory lookup, multimodal setup, `SharedComponents` creation,
and `retry_config` into a private helper, then have both `new` and `new_epd`
call it and only განსხვავate the pipeline builders and `router_type` assignment.

@mergify

mergify Bot commented Jul 7, 2026

Copy link
Copy Markdown
Contributor

Hi @chenht2022, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch:

git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease

@mergify mergify Bot added the needs-rebase PR has merge conflicts that need to be resolved label Jul 7, 2026
Route multimodal items through TokenSpeed encode workers before prefill.

Add the encoder gRPC proto/client/servicer wiring and thread encode bootstrap metadata through gateway request execution.

Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
Signed-off-by: chenht2022 <chenht2022@gmail.com>
@slin1237
slin1237 force-pushed the epd-routing-clean branch from 83db4b2 to e535601 Compare July 7, 2026 13:25
@mergify mergify Bot removed the needs-rebase PR has merge conflicts that need to be resolved label Jul 7, 2026
Signed-off-by: chenht2022 <chenht2022@gmail.com>
@slin1237
slin1237 force-pushed the epd-routing-clean branch from e535601 to cfecdc8 Compare July 7, 2026 13:32

@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: cfecdc88c2

ℹ️ 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 +302 to +303
if Self::matches_selector(pod, &config.encode_selector) {
Some(PodType::Encode)

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 Guard empty encode selector in PD discovery

In plain PD service discovery encode_selector is left empty, but matches_selector returns true for an empty selector because all() over no labels succeeds. Since this check now runs before the prefill/decode checks, every labeled prefill/decode pod is classified as PodType::Encode, so handle_pod_event registers no Prefill/Decode workers and PD routing/readiness never sees a usable pair. Check that the encode selector is non-empty before matching it.

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 python-bindings Python bindings changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants