Repository navigation
Optimize TTFT in PD Disaggregation: Decouple prefill response and decode streaming - #1150
usernamehaha2022 wants to merge 3 commits into
Conversation
📝 WalkthroughWalkthroughModified the PD router to optimize streaming requests without logprobs by spawning a background task for the prefill request instead of awaiting it concurrently with the decode request. The main path now awaits only the decode request and returns immediately, while preserving original behavior for non-streaming or logprob-returning cases. Changes
Sequence Diagram(s)sequenceDiagram
participant Router
participant PrefillTask as Prefill Task<br/>(Background)
participant DecodeReq as Decode Request
participant Client
rect rgba(76, 175, 80, 0.5)
Note over Router,Client: Fast-Path (Streaming, No Logprobs)
Router->>Router: Emit RequestReceivedEvent
Router->>+PrefillTask: Spawn background task
PrefillTask->>PrefillTask: Send prefill request<br/>(fire-and-forget)
deactivate PrefillTask
Router->>+DecodeReq: Await decode request
DecodeReq-->>-Router: Decode response
Router->>Client: Return streaming response<br/>(None prefill logprobs)
end
rect rgba(33, 150, 243, 0.5)
Note over Router,Client: Original Path (Non-Streaming or With Logprobs)
Router->>Router: Emit RequestReceivedEvent
par
Router->>PrefillTask: Send prefill request
and
Router->>DecodeReq: Send decode request
end
PrefillTask-->>Router: Prefill response
DecodeReq-->>Router: Decode response
Router->>Router: Process both responses<br/>(validate, store logprobs if needed)
Router->>Client: Return response
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
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. Comment |
|
Hi @usernamehaha2022, the DCO sign-off check has failed. All commits must include a 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-leaseTo sign off future commits automatically:
|
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: b24b28508b
ℹ️ 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".
| tokio::spawn(async move { | ||
| match prefill_request.send().await { | ||
| Ok(res) => { | ||
| let status = res.status(); | ||
| if !status.is_success() { |
There was a problem hiding this comment.
Cancel fire-and-forget prefill on decode transport failure
The fast path detaches prefill_request.send() in a background task without keeping a handle, so when decode_request.send() fails the function returns while prefill keeps running unmanaged. In PD mode, prefill can block waiting for decode bootstrap (the same hang risk called out by the try_join! comment in the non-fast path), so repeated decode transport errors can strand prefill requests until timeout and consume worker/connection capacity across retries. Please keep a cancellation path (or avoid detached spawn) for prefill when decode fails before streaming starts.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Code Review
This pull request introduces a Time To First Token (TTFT) optimization for streaming requests by backgrounding prefill requests when logprobs are not required. Feedback highlights that this fire-and-forget approach loses the immediate cancellation benefits of the previous tokio::try_join! implementation, which could lead to hanging requests if the prefill fails. Additionally, it was recommended to log errors when consuming the backgrounded prefill response body rather than ignoring them.
| if context.is_stream && !context.return_logprob { | ||
| // Send prefill request in background — don't block decode's stream on it | ||
| let prefill_url_for_log = prefill.url().to_string(); | ||
| #[expect( | ||
| clippy::disallowed_methods, | ||
| reason = "fire-and-forget prefill; gateway shutdown need not wait for prefill response" | ||
| )] | ||
| tokio::spawn(async move { | ||
| match prefill_request.send().await { | ||
| Ok(res) => { | ||
| let status = res.status(); | ||
| if !status.is_success() { | ||
| error!( | ||
| "Prefill server returned error (background) prefill_url={} status={}", | ||
| prefill_url_for_log, status | ||
| ); | ||
| } | ||
| // Consume response body to release the HTTP connection back to pool | ||
| let _ = res.bytes().await; | ||
| } | ||
| Err(e) => { | ||
| error!( | ||
| "Prefill request failed (background) prefill_url={} error={}", | ||
| prefill_url_for_log, e | ||
| ); | ||
| } | ||
| } | ||
| }); | ||
|
|
||
| events::RequestReceivedEvent {}.emit(); | ||
| // Only await decode response — start proxying immediately | ||
| let decode_result = decode_request.send().await; | ||
|
|
||
| let (prefill_response, decode_response) = match pd_result { | ||
| Ok((prefill_resp, decode_resp)) => (prefill_resp, decode_resp), | ||
| Err(e) => { | ||
| error!("PD request transport error, both sides aborted: {e}"); | ||
| // Don't record_outcome here — the caller (execute_dual_dispatch) | ||
| // records outcomes from the response status after we return. | ||
| return error::bad_gateway( | ||
| "PD disaggregation request failed", | ||
| format!("Transport error: {e}"), | ||
| ); | ||
| events::RequestReceivedEvent {}.emit(); | ||
|
|
||
| match decode_result { | ||
| Ok(res) => { | ||
| let status = StatusCode::from_u16(res.status().as_u16()) | ||
| .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); | ||
| debug!("Decode response status (fast path): {}", status); | ||
|
|
||
| if !status.is_success() { | ||
| error!( | ||
| "Decode server returned error decode_url={} status={}", | ||
| decode.url(), | ||
| status | ||
| ); | ||
| return self | ||
| .handle_decode_error_response(res, &context, prefill, decode) | ||
| .await; | ||
| } | ||
|
|
||
| let response_headers = header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| self.create_streaming_response( | ||
| res.bytes_stream(), | ||
| status, | ||
| None, | ||
| false, | ||
| None, | ||
| Some(response_headers), | ||
| prefill, | ||
| decode, | ||
| ) | ||
| } | ||
| Err(e) => { | ||
| error!( | ||
| decode_url = %decode.url(), | ||
| error = %e, | ||
| "Decode request failed" | ||
| ); | ||
| error::bad_gateway("decode_server_error", format!("Decode server error: {e}")) | ||
| } | ||
| } |
There was a problem hiding this comment.
This optimization for the happy path introduces a potential performance degradation and correctness issue in failure scenarios. In the original implementation, tokio::try_join! ensures that if the prefill request fails, the decode request is canceled immediately. In this new fast path, the prefill request is sent in a fire-and-forget tokio::spawn task. If the prefill request fails, the decode request will likely hang until it times out. This violates the principle of graceful error handling for external interactions, as a failure in the background task leads to a poor user experience.
References
- Do not introduce unhandled failure states in code that interacts with external systems. Handle errors gracefully to avoid hanging requests or inconsistent states.
| ); | ||
| } | ||
| // Consume response body to release the HTTP connection back to pool | ||
| let _ = res.bytes().await; |
There was a problem hiding this comment.
The error from res.bytes().await is ignored. According to repository guidelines, potential failures should be logged as warnings to aid in debugging rather than being silently ignored. This ensures that issues like connection pool exhaustion or network errors are visible.
| let _ = res.bytes().await; | |
| if let Err(e) = res.bytes().await { | |
| warn!( | |
| "Failed to consume prefill response body (background) prefill_url={} error={}", | |
| prefill_url_for_log, e | |
| ); | |
| } |
References
- Instead of silently ignoring potential failures, log them as warnings to aid in debugging.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/routers/http/pd_router.rs`:
- Around line 639-710: The background prefill spawn swallows transport/non-2xx
errors and never signals them back, letting a failed prefill become a successful
stream; fix by creating a small tokio::sync::mpsc::channel(1) (e.g., let
(prefill_tx, mut prefill_rx) = mpsc::channel(1)) before spawning, move
prefill_tx into the spawned task and have that task send an Err message (or
error string) on failure and an Ok/empty message on success, then after you
obtain decode_result::Ok(res) and before calling create_streaming_response(...)
poll the channel non-blockingly (prefill_rx.try_recv()) and if you see a prefill
failure return the appropriate error path (e.g., call
handle_decode_error_response or error::bad_gateway) instead of proceeding to
create_streaming_response; update references around prefill_request/send, the
spawned task, prefill_url_for_log, create_streaming_response, and
handle_decode_error_response accordingly.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 91f6e7e7-b76f-465d-ba61-fb3415cb2278
📒 Files selected for processing (1)
model_gateway/src/routers/http/pd_router.rs
| if context.is_stream && !context.return_logprob { | ||
| // Send prefill request in background — don't block decode's stream on it | ||
| let prefill_url_for_log = prefill.url().to_string(); | ||
| #[expect( | ||
| clippy::disallowed_methods, | ||
| reason = "fire-and-forget prefill; gateway shutdown need not wait for prefill response" | ||
| )] | ||
| tokio::spawn(async move { | ||
| match prefill_request.send().await { | ||
| Ok(res) => { | ||
| let status = res.status(); | ||
| if !status.is_success() { | ||
| error!( | ||
| "Prefill server returned error (background) prefill_url={} status={}", | ||
| prefill_url_for_log, status | ||
| ); | ||
| } | ||
| // Consume response body to release the HTTP connection back to pool | ||
| let _ = res.bytes().await; | ||
| } | ||
| Err(e) => { | ||
| error!( | ||
| "Prefill request failed (background) prefill_url={} error={}", | ||
| prefill_url_for_log, e | ||
| ); | ||
| } | ||
| } | ||
| }); | ||
|
|
||
| events::RequestReceivedEvent {}.emit(); | ||
| // Only await decode response — start proxying immediately | ||
| let decode_result = decode_request.send().await; | ||
|
|
||
| let (prefill_response, decode_response) = match pd_result { | ||
| Ok((prefill_resp, decode_resp)) => (prefill_resp, decode_resp), | ||
| Err(e) => { | ||
| error!("PD request transport error, both sides aborted: {e}"); | ||
| // Don't record_outcome here — the caller (execute_dual_dispatch) | ||
| // records outcomes from the response status after we return. | ||
| return error::bad_gateway( | ||
| "PD disaggregation request failed", | ||
| format!("Transport error: {e}"), | ||
| ); | ||
| events::RequestReceivedEvent {}.emit(); | ||
|
|
||
| match decode_result { | ||
| Ok(res) => { | ||
| let status = StatusCode::from_u16(res.status().as_u16()) | ||
| .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); | ||
| debug!("Decode response status (fast path): {}", status); | ||
|
|
||
| if !status.is_success() { | ||
| error!( | ||
| "Decode server returned error decode_url={} status={}", | ||
| decode.url(), | ||
| status | ||
| ); | ||
| return self | ||
| .handle_decode_error_response(res, &context, prefill, decode) | ||
| .await; | ||
| } | ||
|
|
||
| let response_headers = header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| self.create_streaming_response( | ||
| res.bytes_stream(), | ||
| status, | ||
| None, | ||
| false, | ||
| None, | ||
| Some(response_headers), | ||
| prefill, | ||
| decode, | ||
| ) | ||
| } | ||
| Err(e) => { | ||
| error!( | ||
| decode_url = %decode.url(), | ||
| error = %e, | ||
| "Decode request failed" | ||
| ); | ||
| error::bad_gateway("decode_server_error", format!("Decode server error: {e}")) | ||
| } |
There was a problem hiding this comment.
Don't let prefill failures turn into a successful stream.
This branch never surfaces the prefill result back to the caller: transport errors and non-2xx responses are only logged in the spawned task. That means a failed prefill can now become a 200 streaming response, skip retry, and later be recorded as a successful outcome for the prefill worker because execute_dual_dispatch only sees the decode-side status. If decode is waiting on a room that prefill never bootstrapped, the client gets a hung stream instead of the previous retriable error. Please preserve a failure signal from the background prefill path before committing the response, or at least feed prefill failures back into retry/outcome handling.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/http/pd_router.rs` around lines 639 - 710, The
background prefill spawn swallows transport/non-2xx errors and never signals
them back, letting a failed prefill become a successful stream; fix by creating
a small tokio::sync::mpsc::channel(1) (e.g., let (prefill_tx, mut prefill_rx) =
mpsc::channel(1)) before spawning, move prefill_tx into the spawned task and
have that task send an Err message (or error string) on failure and an Ok/empty
message on success, then after you obtain decode_result::Ok(res) and before
calling create_streaming_response(...) poll the channel non-blockingly
(prefill_rx.try_recv()) and if you see a prefill failure return the appropriate
error path (e.g., call handle_decode_error_response or error::bad_gateway)
instead of proceeding to create_streaming_response; update references around
prefill_request/send, the spawned task, prefill_url_for_log,
create_streaming_response, and handle_decode_error_response accordingly.
|
@usernamehaha2022 thanks for contributing. Could you fix the failed CI checks? |
|
This pull request has been automatically marked as stale because it has not had any activity within 14 days. It will be automatically closed if no further activity occurs within 16 days. Leave a comment if you feel this pull request should remain open. Thank you! |
|
This pull request has been automatically closed due to inactivity. Please feel free to reopen if you intend to continue working on it. Thank you! |
Description
Problem
This PR improves the TTFT of PD HTTP routing in the common streaming path.
In the previous implementation of
execute_dual_dispatch_internal, the router waited for both prefill and decode HTTP responses before starting to proxy the decode stream. This adds avoidable latency to the first streamed token for the common case where:Since input logprob merging is only needed when return_logprob=true, waiting for the prefill HTTP response in the no-logprob streaming path does not provide value on the critical path, but does delay stream startup.
The goal of this PR is to reduce that overhead and improve end-to-end streaming responsiveness in PD mode.
Solution
This PR updates the PD HTTP router in pd_router.rs:
decode streaming is proxied to the client as soon as the decode side is ready.
process_prefill_response()and the original dual-response path when return_logprob=true.In short, this change narrows the optimization to the most common latency-sensitive streaming scenario while keeping the logprob path and non-streaming path behavior unchanged.
Test Plan
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit