feat(messages api): support streaming - #280
Conversation
Summary of ChangesHello @key4ng, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request introduces a comprehensive pipeline architecture for the Anthropic router, enhancing its functionality, security, and observability. The changes include a structured approach to processing Messages API requests, improved streaming support with detailed metrics, and security measures to prevent credential leakage and DoS attacks. The new pipeline architecture provides a more modular and maintainable design, allowing for easier addition of new features and improvements in the future. Highlights
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
There was a problem hiding this comment.
Code Review
This pull request introduces a significant refactoring of the Anthropic router, replacing separate streaming and non-streaming handlers with a unified, stage-based pipeline architecture, which improves maintainability and extensibility. The current implementation has critical reliability issues, including flawed worker load management leading to inaccurate routing decisions, and a data corruption bug in MetricsStream that can cause duplicated data in streaming responses. A critical security vulnerability related to unbounded memory consumption in the SSE parser has been identified, but due to its cross-cutting nature and design implications, it should be addressed in a dedicated pull request. While there are excellent security enhancements in models.rs preventing credential leakage and adding response size limits, the remaining critical issues, along with a minor documentation issue, need to be addressed before merging.
554941f to
79636aa
Compare
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces streaming support for the Anthropic Messages API by refactoring request handling into a well-structured, multi-stage pipeline, enhancing modularity, maintainability, and extensibility. It also adds detailed streaming metrics. Several minor cleanups and improvements are suggested, such as removing unused code, eliminating leftover helper methods, and ensuring consistent metrics calculation for both streaming and non-streaming responses.
cee7677 to
8156b27
Compare
📝 WalkthroughWalkthroughAdds a staged Messages pipeline for the Anthropic router: introduces RequestContext and SharedComponents, a MessagesPipeline wiring WorkerSelection, RequestBuilding, RequestExecution, and ResponseProcessing stages; removes legacy streaming/non-streaming handlers; adds response size limits, Anthropic worker filtering, header propagation rules, and request validation. Changes
Sequence DiagramsequenceDiagram
participant Client
participant Router as AnthropicRouter
participant Pipeline as MessagesPipeline
participant WS as WorkerSelectionStage
participant RB as RequestBuildingStage
participant RE as RequestExecutionStage
participant RP as ResponseProcessingStage
participant Registry as WorkerRegistry
participant Worker as RemoteWorker
Client->>Router: route_messages(request, headers, model_id)
Router->>Pipeline: execute(request, headers, model_id)
Pipeline->>WS: execute(ctx)
WS->>Registry: find_best_worker_for_model(model_id)
Registry-->>WS: Arc<Worker>
WS-->>Pipeline: Ok(None)
Pipeline->>RB: execute(ctx)
RB->>RB: build URL and propagate headers
RB-->>Pipeline: Ok(None)
Pipeline->>RE: execute(ctx)
RE->>Worker: POST /v1/messages (JSON + headers) with timeout
alt worker responds
Worker-->>RE: Response
RE-->>Pipeline: Ok(None)
else error / timeout
RE-->>Pipeline: Err(Response)
end
Pipeline->>RP: execute(ctx)
alt streaming
RP->>RP: wrap stream with LoadTrackingStream
RP-->>Pipeline: Ok(Some(streaming Response))
else non-streaming
RP->>RP: read/parse body, enforce limits, record metrics
RP-->>Pipeline: Ok(Some(Response))
end
Pipeline-->>Router: Response
Router-->>Client: Response
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes 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)
Comment |
|
@coderabbitai review |
✅ Actions performedReview triggered.
|
There was a problem hiding this comment.
Actionable comments posted: 5
🤖 Fix all issues with AI agents
In `@model_gateway/src/routers/anthropic/models.rs`:
- Around line 73-92: The code currently calls response.bytes().await (via
req_builder.send().await -> response) and only then checks bytes.len(), which
can buffer an unbounded body; change this to first inspect
response.content_length() and immediately return the error if Some(len) >
MAX_RESPONSE_SIZE, and if content_length is unknown use response.bytes_stream()
to consume the body incrementally while summing chunk lengths and aborting as
soon as the accumulated size exceeds MAX_RESPONSE_SIZE; update the code paths
around req_builder.send().await, response.content_length(), and
response.bytes_stream() so the body is never fully buffered beyond the
configured MAX_RESPONSE_SIZE.
In `@model_gateway/src/routers/anthropic/stages/request_execution.rs`:
- Around line 58-95: The circuit breaker is only recorded in
request_execution.rs via worker.record_outcome(status.is_success()), which
misreports successes when later parsing or streaming fails; move or defer
recording success into ResponseProcessingStage so the outcome is recorded only
after successful parsing/consumption, and add explicit failure recordings: call
worker.record_outcome(false) inside handle_success_response when JSON parsing
fails, and inside LoadTrackingStream::drop when the stream is interrupted
(ensure you can detect interruption vs successful completion), and remove or
guard the existing record_outcome call in request_execution.rs to avoid
double-reporting.
- Around line 17-68: The request builder currently calls
.timeout(Duration::from_secs(DEFAULT_WORKER_TIMEOUT_SECS)) which hard-codes 120s
and overrides the gateway/client timeout; remove the hard-coded
DEFAULT_WORKER_TIMEOUT_SECS usage in RequestExecutionStage::execute and instead
obtain the configured timeout from the shared configuration (e.g.
SharedComponents or request_timeout_secs) — either by adding a timeout field to
RequestExecutionStage via new(http_client, request_timeout_secs) or by reading
ctx.shared_components.request_timeout_secs inside execute — then apply that
value (converted to Duration) to the reqwest request builder (or omit
per-request timeout to rely on the client-level timeout) so gateway-wide timeout
configuration is respected.
In `@model_gateway/src/routers/anthropic/stages/response_processing.rs`:
- Around line 280-346: The handler handle_error_response currently calls
response.bytes().await which buffers the entire body before checking
MAX_ERROR_RESPONSE_SIZE; change this to consume response.bytes_stream() and read
chunks up to MAX_ERROR_RESPONSE_SIZE (accumulating into a Vec<u8> and stopping
when the limit is exceeded), returning a truncated/overflow message if exceeded,
and handling stream errors similarly to the Err branch; ensure you still produce
a lossily-decoded body_preview, log the body_size and overflow event, and
preserve the existing metrics, worker.decrement_load() and final Err((status,
body).into_response()) behavior.
In `@model_gateway/src/routers/anthropic/stages/worker_selection.rs`:
- Around line 28-74: The find_best_worker_for_model function currently calls
worker_registry.get_workers_filtered with None provider and thus can pick
wildcard workers across providers; change it to be provider-aware like
models.rs: detect if multiple providers exist in the cluster and, when Anthropic
credentials may be present for the request, pass Some(ProviderType::Anthropic)
(or the appropriate ProviderType) into get_workers_filtered instead of None;
update the call site in find_best_worker_for_model (and keep using
supports_model and min_by_key(load) as before) so only workers from the intended
provider are considered and credentials cannot leak to other providers.
🧹 Nitpick comments (2)
model_gateway/src/routers/anthropic/stages/response_processing.rs (1)
225-278: Return Ok(Some) for successful responses to match StageResult contract.
Err(...)is documented as the error path; consider usingOk(Some(...))for success (and making the pipeline log message more neutral if you adopt this).♻️ Suggested adjustment
- Err((StatusCode::OK, Json(message)).into_response()) + Ok(Some((StatusCode::OK, Json(message)).into_response()))model_gateway/src/routers/anthropic/context.rs (1)
120-155: Prefer validated streaming flag when available.
If the validation stage ever normalizes or overrides streaming,is_streaming()could diverge from the pipeline state. Consider consultingstate.validationfirst and falling back to the request flag.♻️ Suggested adjustment
pub fn is_streaming(&self) -> bool { - self.input.request.stream.unwrap_or(false) + self.state + .validation + .as_ref() + .map(|v| v.is_streaming) + .unwrap_or_else(|| self.input.request.stream.unwrap_or(false)) }
ea2628d to
39bc681
Compare
346ac2b to
acfe956
Compare
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 Fix all issues with AI agents
In `@model_gateway/src/routers/anthropic/stages/response_processing.rs`:
- Around line 168-176: The Content-Type check for SSE is case-sensitive and can
miss mixed-case headers; update the check in the response handling (where
content_type is derived from response.headers().get(header::CONTENT_TYPE)) to
perform a case-insensitive comparison — e.g., convert content_type to lowercase
with to_ascii_lowercase() (or use a case-insensitive comparison) before calling
contains("text/event-stream") so the if condition that detects SSE streams
correctly matches values like "Text/Event-Stream" or "text/Event-Stream".
In `@model_gateway/src/routers/anthropic/utils.rs`:
- Around line 115-142: The function filter_by_anthropic_provider currently
treats missing provider (default_provider() -> None) as compatible with a single
explicit provider which can leak Anthropic credentials; change the
provider-detection logic to treat None as a distinct provider value so any mix
of Some(...) and None counts as multiple providers. Concretely, in
filter_by_anthropic_provider update the tracking variable (first_provider) to
hold Option<Option<ProviderType>> and compare the full Option<ProviderType>
returned by Worker::default_provider() when scanning workers; if you detect
differing Option<ProviderType> values set has_multiple_providers true as before,
and keep the existing filtering branch that only keeps workers where
default_provider() == Some(ProviderType::Anthropic).
- Around line 85-108: The current read_response_body_limited function decodes
each chunk with from_utf8_lossy which can corrupt multibyte UTF-8 sequences
split across chunks; fix it by accumulating raw bytes into a Vec<u8> (instead of
appending chunk-decoded Strings), enforcing the max_size cap by checking
total_size after adding each chunk and returning ReadBodyResult::TooLarge if
exceeded, and only after the stream completes decode the entire buffer once
using String::from_utf8 (or map the UTF-8 error to ReadBodyResult::Error) and
return ReadBodyResult::Ok(body) on success; keep the existing chunk read error
handling (Err branch) returning ReadBodyResult::Error(e.to_string()) and update
references within read_response_body_limited accordingly.
🧹 Nitpick comments (1)
model_gateway/src/routers/anthropic/stages/request_execution.rs (1)
37-118: Emit router metrics on dispatch failures.When
send()fails,ResponseProcessingStagenever runs, so request/duration/error metrics may be missing for timeouts/connect errors. Consider recording them here for observability parity.🧭 Suggested metrics hook for send failures
@@ - let url = &http_request.url; + let url = &http_request.url; + let model_id = &ctx.input.model_id; + let start_time = ctx.start_time; + let is_streaming = ctx.is_streaming(); @@ // Record circuit breaker failure worker.record_outcome(false); + + // Record router metrics for dispatch failures (response_processing won't run) + Metrics::record_router_request( + metrics_labels::ROUTER_HTTP, + metrics_labels::BACKEND_EXTERNAL, + metrics_labels::CONNECTION_HTTP, + model_id, + "messages", + bool_to_static_str(is_streaming), + ); + Metrics::record_router_duration( + metrics_labels::ROUTER_HTTP, + metrics_labels::BACKEND_EXTERNAL, + metrics_labels::CONNECTION_HTTP, + model_id, + "messages", + start_time.elapsed(), + ); + Metrics::record_router_error( + metrics_labels::ROUTER_HTTP, + metrics_labels::BACKEND_EXTERNAL, + metrics_labels::CONNECTION_HTTP, + model_id, + "messages", + metrics_labels::ERROR_BACKEND, + );➕ Import metrics utilities
-use crate::routers::{anthropic::context::RequestContext, error}; +use crate::{ + observability::metrics::{bool_to_static_str, metrics_labels, Metrics}, + routers::{anthropic::context::RequestContext, error}, +};
Signed-off-by: ppraneth <pranethparuchuri@gmail.com>
Description
Problem
The Anthropic Messages API (
/v1/messages) lacked streaming support, limiting real-time response delivery for clients using thestream: trueparameter.Solution
Implement full SSE (Server-Sent Events) streaming for the Anthropic Messages API, following the Anthropic streaming specification with proper event types and delta handling.
Changes
message_start,content_block_start,content_block_delta,message_delta,message_stop)validation- Request validationworker_selection- Route to appropriate backendrequest_building- Transform request for backendrequest_execution- Execute request with streaming supportresponse_processing- Handle streaming SSE responsedispatch_metadata- Track request metadatatext,thinking,tool_useTest Plan
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Improvements
Chores