feat(realtime-api): realtime websocket handler - #637
Conversation
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request significantly enhances the Realtime API gateway by introducing a WebSocket transport layer. Previously, the gateway only supported REST endpoints for token generation, lacking the necessary bidirectional streaming for interactive realtime sessions. The new implementation provides a dedicated WebSocket endpoint that upgrades client connections, intelligently routes them to appropriate upstream workers, and transparently proxies messages, thereby enabling full server-to-server realtime communication. Highlights
Changelog
Activity
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
|
📝 WalkthroughWalkthroughAdds workspace deps and WebSocket realtime support: new bidirectional TLS WebSocket proxy, a /v1/realtime WebSocket handler that registers sessions with a realtime registry, registry reaper integration in server startup, and registry API changes (dynamic sizing, get/set state). Changes
Sequence DiagramsequenceDiagram
participant Client
participant WsHandler as "ws_handler (HTTP)"
participant Registry as "RealtimeRegistry"
participant Proxy as "WebSocket Proxy"
participant Upstream as "Upstream Realtime"
Client->>WsHandler: GET /v1/realtime?model=...
WsHandler->>WsHandler: Authenticate & select worker
WsHandler->>Registry: Register session -> cancel_token
WsHandler->>WsHandler: Build upstream WebSocket URL
WsHandler->>Client: Upgrade to WebSocket
WsHandler->>Proxy: hand off WebSocket + session info
Proxy->>Upstream: Establish TLS WebSocket connection
Proxy->>Registry: set session state = Connected
par Bidirectional forwarding
loop client -> upstream
Client->>Proxy: WS message (Text/Binary/Ping/Close)
Proxy->>Proxy: parse/log ClientEvent
Proxy->>Upstream: forward message
end
and
loop upstream -> client
Upstream->>Proxy: WS message
Proxy->>Proxy: parse/log ServerEvent
Proxy->>Client: forward message
end
end
Client-->>Proxy: Close or cancel_token triggered
Proxy->>Upstream: Close upstream connection
Proxy->>Registry: set session state = Disconnected
Estimated Code Review Effort🎯 4 (Complex) | ⏱️ ~50 minutes Possibly Related Issues
Possibly Related PRs
Suggested Reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces a WebSocket handler for the OpenAI Realtime API, utilizing good patterns for stream management and session handling. However, a critical Resource Exhaustion (DoS) vulnerability was identified: sessions are registered before the WebSocket upgrade is completed, potentially filling the registry's capacity with stale 'Pending' sessions. Beyond this, consider optimizing performance by avoiding unnecessary data clones in the message forwarding hot path, and enhancing robustness through URL encoding for query parameters and improved authentication header validation.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 3426ebe914
ℹ️ 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".
There was a problem hiding this comment.
Actionable comments posted: 5
🤖 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/openai/realtime/proxy.rs`:
- Around line 66-103: When one forwarder join handle finishes in the
tokio::select! (client_to_upstream or upstream_to_client), explicitly abort the
sibling join handle (call .abort() on the other JoinHandle), then await that
handle to ensure it has terminated and log any JoinError; do this after
signalling cancel_token.cancel() so both forward_client_to_upstream and
forward_upstream_to_client are deterministically stopped and cleaned up instead
of being dropped and left running.
- Around line 47-49: Wrap the call to
tokio_tungstenite::connect_async_tls_with_config(request, None, false,
Some(connector)).await in a tokio::time::timeout(...) with the standard
connection timeout duration used elsewhere, and handle a timeout by mapping the
timeout error into the same error type via .map_err(...) so the dial failure
returns a clear timeout error instead of hanging; keep the existing variable
names (request, connector, upstream_ws) and ensure the map_err path produces the
same error shape used by the surrounding function.
In `@model_gateway/src/routers/openai/realtime/ws.rs`:
- Around line 115-117: The build_upstream_ws_url function currently interpolates
model directly into the query string which allows reserved chars to break URL
semantics; update build_upstream_ws_url(worker_url: &str, model: &str) to
percent-encode the model value before insertion (e.g., using the
percent-encoding or url crate / Url::parse_with_params) and then format the
final URL from the trimmed base and the encoded model so characters like &, =, #
are safely escaped in the query component.
- Around line 63-66: The code currently coerces non-UTF8 Authorization headers
into an empty string (auth_str) which then becomes a valid empty header
upstream; instead, in the websocket handler where extract_auth_header(...) and
auth_header_value are used, treat a to_str() failure as an invalid header and
reject the request with a 401 rather than converting to "". Replace the
unwrap_or("") behavior for auth_header_value -> auth_str so that when
HeaderValue::to_str() returns Err you return an immediate 401 Unauthorized
(mirroring rest.rs behavior) and only forward a valid UTF-8 auth_str upstream;
ensure the error path references extract_auth_header, auth_header_value,
auth_str and worker.api_key() to locate the change.
In `@model_gateway/src/server.rs`:
- Around line 566-570: start_reaper currently spawns an unguarded background
task (called from realtime_registry.start_reaper) every time build_app/test
helpers run, causing duplicate reapers and losing the CancellationToken; make it
idempotent by adding an internal guard or stored token: either use a Once or
OnceLock inside realtime_registry to run start_reaper only once, or have
start_reaper store and return a CancellationToken field on the registry (check
if it's already Some and return early) so repeated calls do nothing and the
token can be used for cleanup; update callers (build_app/test helpers) to rely
on the registry’s idempotent behavior instead of spawning new tasks.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: ca462f6e-d123-425b-9071-15dc3255c104
📒 Files selected for processing (6)
Cargo.tomlmodel_gateway/Cargo.tomlmodel_gateway/src/routers/openai/realtime/mod.rsmodel_gateway/src/routers/openai/realtime/proxy.rsmodel_gateway/src/routers/openai/realtime/ws.rsmodel_gateway/src/server.rs
3426ebe to
1586b8a
Compare
There was a problem hiding this comment.
♻️ Duplicate comments (1)
model_gateway/src/routers/openai/realtime/ws.rs (1)
119-121:⚠️ Potential issue | 🟡 MinorURL-encode
modelbefore interpolating into the query string.Model names containing reserved characters (
&,=,#,?, spaces) would corrupt the URL. Theurlcrate is already available (transitive dependency).🔧 Proposed fix
fn build_upstream_ws_url(worker_url: &str, model: &str) -> String { let base = worker_url.trim_end_matches('/'); - format!("{base}/v1/realtime?model={model}") + let encoded_model: String = + url::form_urlencoded::byte_serialize(model.as_bytes()).collect(); + format!("{base}/v1/realtime?model={encoded_model}") }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/routers/openai/realtime/ws.rs` around lines 119 - 121, The build_upstream_ws_url function currently interpolates model raw into the query string; percent-encode the model value before formatting to avoid breaking URLs when it contains reserved characters. Use the url crate (e.g., url::form_urlencoded::byte_serialize or url::percent_encoding::utf8_percent_encode) to encode the model string, then call format!("{base}/v1/realtime?model={encoded_model}") in build_upstream_ws_url so the query parameter is safely encoded.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@model_gateway/src/routers/openai/realtime/ws.rs`:
- Around line 119-121: The build_upstream_ws_url function currently interpolates
model raw into the query string; percent-encode the model value before
formatting to avoid breaking URLs when it contains reserved characters. Use the
url crate (e.g., url::form_urlencoded::byte_serialize or
url::percent_encoding::utf8_percent_encode) to encode the model string, then
call format!("{base}/v1/realtime?model={encoded_model}") in
build_upstream_ws_url so the query parameter is safely encoded.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 8d073e4a-4d2f-4f64-b039-7a6aa9502f47
📒 Files selected for processing (7)
Cargo.tomlmodel_gateway/Cargo.tomlmodel_gateway/src/routers/openai/realtime/mod.rsmodel_gateway/src/routers/openai/realtime/proxy.rsmodel_gateway/src/routers/openai/realtime/registry.rsmodel_gateway/src/routers/openai/realtime/ws.rsmodel_gateway/src/server.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 1586b8a606
ℹ️ 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".
| error!(session_id, error = %e, "Realtime WebSocket proxy error"); | ||
| } |
There was a problem hiding this comment.
Record WebSocket proxy failures on the selected worker
ws_handler uses select_worker, which filters by circuit_breaker().can_execute(), but this error branch only logs and never calls worker.record_outcome(false). In a deployment where a worker’s realtime endpoint is failing (bad upstream config, repeated handshake/runtime errors), /v1/realtime traffic will keep selecting that worker because realtime failures never feed into the circuit breaker, causing persistent user-facing connection failures instead of tripping and shedding to healthier workers.
Useful? React with 👍 / 👎.
Signed-off-by: yifeliu <yifengliu9@gmail.com>
1586b8a to
cd80e62
Compare
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/openai/realtime/proxy.rs`:
- Around line 87-106: The select block drops the surviving task's JoinHandle
(client_to_upstream or upstream_to_client) after calling cancel_token.cancel(),
detaching it instead of awaiting completion; update the logic so after detecting
one branch finished you call await on the other task's JoinHandle (the surviving
client_to_upstream or upstream_to_client handle) to ensure deterministic cleanup
once cancel_token.cancel() is issued, while still keeping the
cancel_token.cancelled() branch; locate the select around those symbols and add
a follow-up await for the other handle (handling/ignoring its Result) after the
select completes to avoid leaving the task detached.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: fdd747dc-8df0-4ebd-bfcb-af79480bc36b
📒 Files selected for processing (7)
Cargo.tomlmodel_gateway/Cargo.tomlmodel_gateway/src/routers/openai/realtime/mod.rsmodel_gateway/src/routers/openai/realtime/proxy.rsmodel_gateway/src/routers/openai/realtime/registry.rsmodel_gateway/src/routers/openai/realtime/ws.rsmodel_gateway/src/server.rs
| // Wait for either task to finish (or cancellation) | ||
| tokio::select! { | ||
| result = client_to_upstream => { | ||
| cancel_token.cancel(); | ||
| debug!(session_id, "Client→upstream task ended"); | ||
| if let Err(e) = result { | ||
| error!(session_id, error = %e, "Client→upstream task panicked"); | ||
| } | ||
| } | ||
| result = upstream_to_client => { | ||
| cancel_token.cancel(); | ||
| debug!(session_id, "Upstream→client task ended"); | ||
| if let Err(e) = result { | ||
| error!(session_id, error = %e, "Upstream→client task panicked"); | ||
| } | ||
| } | ||
| () = cancel_token.cancelled() => { | ||
| debug!(session_id, "Session cancelled via token"); | ||
| } | ||
| } |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Consider awaiting the surviving task for deterministic cleanup.
When one forwarding task completes, the other's JoinHandle is dropped without awaiting. While cancel_token.cancel() signals the surviving task to exit, dropping the handle detaches it rather than ensuring completion.
Per the past discussion, the author notes that the surviving task exits within one event loop tick via its cancellation branch, making this acceptable in practice. If you want structural guarantees, you could await the surviving handle after cancellation:
// After the select! block, in each branch:
let _ = other_handle.await;This is low-severity given the task's short lifetime and fast cancellation response.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/openai/realtime/proxy.rs` around lines 87 - 106,
The select block drops the surviving task's JoinHandle (client_to_upstream or
upstream_to_client) after calling cancel_token.cancel(), detaching it instead of
awaiting completion; update the logic so after detecting one branch finished you
call await on the other task's JoinHandle (the surviving client_to_upstream or
upstream_to_client handle) to ensure deterministic cleanup once
cancel_token.cancel() is issued, while still keeping the
cancel_token.cancelled() branch; locate the select around those symbols and add
a follow-up await for the other handle (handling/ignoring its Result) after the
select completes to avoid leaving the task detached.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: cd80e62da4
ℹ️ 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".
| )); | ||
|
|
||
| let realtime_routes = Router::new() | ||
| .route("/v1/realtime", get(realtime_ws::ws_handler)) |
There was a problem hiding this comment.
Add connection-lifetime limiting for /v1/realtime
Mounting GET /v1/realtime inside realtime_routes means it only gets the existing concurrency_limit_middleware, but that middleware accounts tokens against the HTTP response body lifetime (TokenGuardBody) rather than the upgraded socket lifetime. For WebSocket upgrades, the 101 response body is dropped immediately, so the token is released while the session remains open; in environments that rely on this limiter to cap concurrent load, clients can open many long-lived realtime sockets and bypass the configured concurrency guard.
Useful? React with 👍 / 👎.
proxy.rs: - Replace full ClientEvent/ServerEvent deserialization with lightweight EventTypeOnly struct for event logging. Avoids parsing large boxed variants (SessionConfig 624B, ResponseCreateParams 384B) on every audio frame in the hot path. - Cache TLS ClientConfig via OnceLock instead of rebuilding it (including cloning all webpki root certificates) per connection. - Remove unused openai_protocol::realtime_events import. registry.rs: - Unify SessionEntry and CallEntry (identical except for id field name) into a single ConnectionEntry type. - Extract shared CRUD and reap logic into ConnectionMap, eliminating ~30 lines of duplicated code across session/call methods and the reaper task. Refs: #637 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
proxy.rs: - Use Cow<str> instead of &str for EventTypeOnly::event_type so serde can handle JSON escape sequences (e.g. \u002E) that require allocation, while preserving zero-copy for the common unescaped case. - Change "Safety:" comment to "INVARIANT:" per repo convention. registry.rs: - Replace two-pass collect-then-remove in reap_stale with single-pass DashMap::retain, eliminating intermediate Vec allocation. Refs: #637 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
proxy.rs: - Replace full ClientEvent/ServerEvent deserialization with lightweight EventTypeOnly struct for event logging. Avoids parsing large boxed variants (SessionConfig 624B, ResponseCreateParams 384B) on every audio frame in the hot path. - Cache TLS ClientConfig via OnceLock instead of rebuilding it (including cloning all webpki root certificates) per connection. - Remove unused openai_protocol::realtime_events import. registry.rs: - Unify SessionEntry and CallEntry (identical except for id field name) into a single ConnectionEntry type. - Extract shared CRUD and reap logic into ConnectionMap, eliminating ~30 lines of duplicated code across session/call methods and the reaper task. Refs: #637 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
proxy.rs: - Use Cow<str> instead of &str for EventTypeOnly::event_type so serde can handle JSON escape sequences (e.g. \u002E) that require allocation, while preserving zero-copy for the common unescaped case. - Change "Safety:" comment to "INVARIANT:" per repo convention. registry.rs: - Replace two-pass collect-then-remove in reap_stale with single-pass DashMap::retain, eliminating intermediate Vec allocation. Refs: #637 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
Description
Fix #245
Problem
The Realtime API gateway currently supports REST endpoints for ephemeral token generation (
/v1/realtime/sessions,/v1/realtime/client_secrets,/v1/realtime/transcription_sessions) but lacks the WebSocket transport required for server-to-server bidirectional streaming — the primary transport for realtime audio/text sessions.Solution
Add a bidirectional WebSocket proxy handler at
GET /v1/realtimethat upgrades incoming client connections, selects an upstream worker via model routing, and transparently forwards messages in both directions. The proxy reuses the existingRealtimeRegistryfor session tracking and the same auth/worker-selection logic as the REST handlers.Changes
model_gateway/src/routers/openai/realtime/ws.rs— Axum WebSocket upgrade handler: validatesmodelquery param, selects a worker, registers a session, and delegates to the proxymodel_gateway/src/routers/openai/realtime/proxy.rs— Bidirectional forwarding engine: splits client and upstream WebSocket streams into independentclient→upstreamandupstream→clienttasks coordinated viaCancellationToken, with structured logging by event type (trace for high-frequency audio deltas, info for session lifecycle, debug for everything else). Builds an explicitrustls TLS connector to avoid depending on the process-level
CryptoProvidermodel_gateway/src/routers/openai/realtime/registry.rs— remove capoacity enforcemant since dashmap grows dynamically.model_gateway/src/routers/openai/realtime/mod.rs— Exports the newproxyandwsmodulesmodel_gateway/src/server.rs— WiresGET /v1/realtimeinto the realtime route group (behind auth + concurrency middleware) and starts the session reaper (60 min max age, 1 min sweep interval)Cargo.toml,model_gateway/Cargo.toml— Addstokio-tungsteniteandwebpki-rootsworkspace dependenciesTest Plan
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Behavior Changes