Repository navigation
chore: remove unnecessary result wrapping in pd_router - #939
Conversation
|
Warning Rate limit exceeded
Your organization is not enrolled in usage-based pricing. Contact your admin to enable usage-based pricing to continue reviews beyond the rate limit, or try again in 4 minutes and 32 seconds. ⌛ How to resolve this issue?After the wait time has elapsed, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout. Please see our FAQ for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (1)
📝 WalkthroughWalkthroughReworked dual-dispatch response handling in Changes
Sequence Diagram(s)sequenceDiagram
participant Client
participant PDRouter
participant PrefillSvc
participant DecodeSvc
Client->>PDRouter: send request
PDRouter->>PrefillSvc: dispatch prefill request (async)
PDRouter->>DecodeSvc: dispatch decode request (async)
PrefillSvc-->>PDRouter: reqwest::Response (prefill_response)
DecodeSvc-->>PDRouter: reqwest::Response (decode_response)
alt decode_response transport error
PDRouter->>Client: return 502 Bad Gateway
else decode_response received
alt decode_response status non-success
PDRouter->>PDRouter: handle_decode_error_response(decode_response)
PDRouter->>Client: return error response from decode handler
else decode_response success
PDRouter->>PDRouter: process_prefill_response(prefill_response)
PDRouter->>Client: build and return combined response
end
end
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 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 @lawrence-harmonic, 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.
Code Review
This pull request refactors the PDRouter to simplify the handling of prefill and decode responses by flattening nested match statements and updating helper functions to accept raw responses. The feedback identifies an opportunity to further clean up the code by removing an unnecessary variable alias and consolidating redundant conditional logic when processing prefill responses to improve readability and maintainability.
| let res = decode_response; | ||
| let status = StatusCode::from_u16(res.status().as_u16()) | ||
| .unwrap_or(StatusCode::INTERNAL_SERVER_ERROR); | ||
| debug!("Decode response status: {}", status); | ||
|
|
||
| // Process prefill response | ||
| let prefill_body = if context.return_logprob { | ||
| match self | ||
| .process_prefill_response( | ||
| prefill_result, | ||
| prefill.url(), | ||
| context.return_logprob, | ||
| ) | ||
| .await | ||
| { | ||
| Ok((_, body)) => body, | ||
| Err(error_response) => return error_response, | ||
| } | ||
| } else { | ||
| // Even if we don't need logprobs, we should check prefill status | ||
| match self | ||
| .process_prefill_response(prefill_result, prefill.url(), false) | ||
| .await | ||
| { | ||
| Ok((_, body)) => body, | ||
| Err(error_response) => return error_response, | ||
| } | ||
| }; | ||
|
|
||
| if context.is_stream { | ||
| // Streaming response | ||
| let prefill_logprobs = if context.return_logprob { | ||
| prefill_body | ||
| .as_ref() | ||
| .and_then(|body| serde_json::from_slice::<Value>(body).ok()) | ||
| .and_then(|json| { | ||
| json.pointer("/meta_info/input_token_logprobs").cloned() | ||
| }) | ||
| } else { | ||
| None | ||
| }; | ||
| if !status.is_success() { | ||
| error!( | ||
| "Decode server returned error status decode_url={} status={}", | ||
| decode.url(), | ||
| status | ||
| ); | ||
|
|
||
| let response_headers = header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| self.create_streaming_response( | ||
| res.bytes_stream(), | ||
| status, | ||
| prefill_logprobs, | ||
| context.return_logprob, | ||
| None, | ||
| Some(response_headers), | ||
| prefill, | ||
| decode, | ||
| ) | ||
| } else { | ||
| // Non-streaming response | ||
| if context.return_logprob { | ||
| self.process_non_streaming_response( | ||
| res, | ||
| status, | ||
| context.return_logprob, | ||
| prefill_body, | ||
| ) | ||
| .await | ||
| } else { | ||
| // Direct passthrough when no logprobs needed | ||
| let response_headers = | ||
| header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| match res.bytes().await { | ||
| Ok(decode_body) => { | ||
| let mut response = Response::new(Body::from(decode_body)); | ||
| *response.status_mut() = status; | ||
| *response.headers_mut() = response_headers; | ||
| response | ||
| } | ||
| Err(e) => { | ||
| error!("Failed to read decode response: {}", e); | ||
| error::internal_error( | ||
| "read_response_failed", | ||
| "Failed to read response", | ||
| ) | ||
| } | ||
| } | ||
| return self | ||
| .handle_decode_error_response(res, &context, prefill, decode) | ||
| .await; | ||
| } | ||
|
|
||
| // Process prefill response | ||
| let prefill_body = if context.return_logprob { | ||
| match self | ||
| .process_prefill_response(prefill_response, prefill.url(), context.return_logprob) | ||
| .await | ||
| { | ||
| Ok((_, body)) => body, | ||
| Err(error_response) => return error_response, | ||
| } | ||
| } else { | ||
| // Even if we don't need logprobs, we should check prefill status | ||
| match self | ||
| .process_prefill_response(prefill_response, prefill.url(), false) | ||
| .await | ||
| { | ||
| Ok((_, body)) => body, | ||
| Err(error_response) => return error_response, | ||
| } | ||
| }; | ||
|
|
||
| if context.is_stream { | ||
| // Streaming response | ||
| let prefill_logprobs = if context.return_logprob { | ||
| prefill_body | ||
| .as_ref() | ||
| .and_then(|body| serde_json::from_slice::<Value>(body).ok()) | ||
| .and_then(|json| json.pointer("/meta_info/input_token_logprobs").cloned()) | ||
| } else { | ||
| None | ||
| }; | ||
|
|
||
| let response_headers = header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| self.create_streaming_response( | ||
| res.bytes_stream(), | ||
| status, | ||
| prefill_logprobs, | ||
| context.return_logprob, | ||
| None, | ||
| Some(response_headers), | ||
| prefill, | ||
| decode, | ||
| ) | ||
| } else { | ||
| // Non-streaming response | ||
| if context.return_logprob { | ||
| self.process_non_streaming_response( | ||
| res, | ||
| status, | ||
| context.return_logprob, | ||
| prefill_body, | ||
| ) | ||
| .await | ||
| } else { | ||
| // Direct passthrough when no logprobs needed | ||
| let response_headers = header_utils::preserve_response_headers(res.headers()); | ||
|
|
||
| match res.bytes().await { | ||
| Ok(decode_body) => { | ||
| let mut response = Response::new(Body::from(decode_body)); | ||
| *response.status_mut() = status; | ||
| *response.headers_mut() = response_headers; | ||
| response | ||
| } | ||
| Err(e) => { | ||
| error!("Failed to read decode response: {}", e); | ||
| error::internal_error("read_response_failed", "Failed to read response") | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
The variable res is an unnecessary alias for decode_response. To improve clarity and reduce redundancy, you can remove the alias and use decode_response directly. Additionally, the logic for processing the prefill response contains redundant match statements and an unnecessary if/else block. Refactoring this to call process_prefill_response once with context.return_logprob simplifies the code and avoids duplication, adhering to repository guidelines.
let status = StatusCode::from_u16(decode_response.status().as_u16())
.unwrap_or(StatusCode::INTERNAL_SERVER_ERROR);
debug!("Decode response status: {}", status);
if !status.is_success() {
error!(
"Decode server returned error status decode_url={} status={}",
decode.url(),
status
);
return self
.handle_decode_error_response(decode_response, &context, prefill, decode)
.await;
}
// Process prefill response
let prefill_body = match self
.process_prefill_response(prefill_response, prefill.url(), context.return_logprob)
.await
{
Ok((_, body)) => body,
Err(error_response) => return error_response,
};
if context.is_stream {
// Streaming response
let prefill_logprobs = if context.return_logprob {
prefill_body
.as_ref()
.and_then(|body| serde_json::from_slice::<Value>(body).ok())
.and_then(|json| json.pointer("/meta_info/input_token_logprobs").cloned())
} else {
None
};
let response_headers = header_utils::preserve_response_headers(decode_response.headers());
self.create_streaming_response(
decode_response.bytes_stream(),
status,
prefill_logprobs,
context.return_logprob,
None,
Some(response_headers),
prefill,
decode,
)
} else {
// Non-streaming response
if context.return_logprob {
self.process_non_streaming_response(
decode_response,
status,
context.return_logprob,
prefill_body,
)
.await
} else {
// Direct passthrough when no logprobs needed
let response_headers = header_utils::preserve_response_headers(decode_response.headers());
match decode_response.bytes().await {
Ok(decode_body) => {
let mut response = Response::new(Body::from(decode_body));
*response.status_mut() = status;
*response.headers_mut() = response_headers;
response
}
Err(e) => {
error!("Failed to read decode response: {}", e);
error::internal_error("read_response_failed", "Failed to read response")
}
}
}
}References
- Refactor match statements to avoid duplication. When arms have common logic, use the match to return the differing value and perform the common logic once.
- Prioritize code simplicity and clarity over micro-optimizations, especially when the performance gain is negligible for typical use cases.
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 619-637: The branching around context.return_logprob is
redundant—both branches call self.process_prefill_response(prefill_response,
prefill.url(), <flag>). Replace the if/else with a single call using
context.return_logprob as the flag, await the result, and extract the body or
return the error_response; update the assignment to prefill_body accordingly so
process_prefill_response, prefill_response, prefill.url(), and
context.return_logprob are used only once.
🪄 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: 1190b347-b80f-4002-8f58-5136416db461
📒 Files selected for processing (1)
model_gateway/src/routers/http/pd_router.rs
Signed-off-by: Lawrence Wu <lawrence.wu@harmonic.fun>
09157ab to
56e8a9f
Compare
|
Hi @lawrence-harmonic, 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:
|
Signed-off-by: Lawrence Wu <lawrence.wu@harmonic.fun>
6e0d55c to
3fe9e1d
Compare
small cleanup from #844, we don't need to do
Ok((prefill_resp, decode_resp)) => (Ok(prefill_resp), Ok(decode_resp))since we already gotOkSummary by CodeRabbit
Release Notes