feat(realtime-api): add Realtime API REST handlers and session registry - #406
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
📝 WalkthroughWalkthroughAdds OpenAI Realtime gateway support: workspace deps (tokio-tungstenite, multer), a DashMap-backed RealtimeRegistry with background reaper, REST token/session endpoints and proxy logic, and integration of realtime registry/routes into AppContext and server routing. Changes
Sequence DiagramsequenceDiagram
participant Client
participant Gateway as Gateway (REST Handler)
participant Registry as RealtimeRegistry
participant WorkerMgr as Worker Manager
participant Upstream as Upstream Service
Client->>Gateway: POST /v1/realtime/client_secrets (body with model)
Gateway->>Gateway: extract model & auth
Gateway->>Registry: (optional) read/write session metadata
Gateway->>WorkerMgr: select_worker(model)
WorkerMgr-->>Gateway: Arc<Worker> (or None)
Gateway->>Upstream: forward POST with auth, headers, body
Upstream-->>Gateway: Response (status, headers, body)
Gateway->>Gateway: proxy_response -> build client Response
Gateway-->>Client: Return upstream status + body
Estimated code review effort🎯 4 (Complex) | ⏱️ ~45 minutes Possibly related issues
Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 3 | ❌ 2❌ Failed checks (2 warnings)
✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
Summary of ChangesHello @pallasathena92, 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 extends the gateway's capabilities by introducing foundational support for OpenAI's Realtime API. It enables clients to generate ephemeral tokens for real-time interactions and establishes an in-memory system for tracking active real-time connections. The changes lay the groundwork for future WebSocket proxy, WebRTC signaling, and MCP interception features, ensuring the gateway can effectively manage and route real-time traffic. 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
|
There was a problem hiding this comment.
Code Review
This pull request introduces support for the OpenAI Realtime API, including REST handlers for ephemeral token generation and an in-memory registry for tracking connections. While the implementation is generally well-structured, there are critical security and reliability concerns: new routes bypass global concurrency limits, and REST handlers lack proper worker load tracking and circuit breaker integration. The session registry reaper also has a potential resource leak for 'Connected' sessions. Further improvements are needed for maintainability and efficiency, specifically addressing code duplication, hardcoded values, and optimizing data access and cleanup in the registry. An unused dependency should also be removed, and Rust use statement conventions should be followed.
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Fix all issues with AI agents
In `@model_gateway/src/routers/openai/realtime/registry.rs`:
- Around line 91-116: Existing session/call entries' CancellationToken objects
are never triggered when entries are removed or replaced, causing awaiting tasks
to leak; update register_session and register_call to check for an existing
entry returned by the map insert/replace and call
existing_entry.cancel_token.cancel() before replacing it, and update
remove_session and remove_call (and the reaper eviction logic) to call
entry.cancel_token.cancel() before dropping/removing the entry so any waiting
tasks are signaled; reference the SessionEntry/CallEntry structs' cancel_token
fields and the register_session/register_call/remove_session/remove_call
functions and the reaper eviction code to apply this fix.
- Around line 83-90: The capacity checks for sessions and calls are vulnerable
to TOCTOU races; change register_session (referencing self.sessions and
self.max_sessions) and register_call (referencing self.calls and self.max_calls)
to use an atomic reservation approach instead of checking len() then inserting.
Add an AtomicUsize (e.g., session_count and call_count) or a semaphore to
perform an atomic fetch_add/try_acquire before inserting, only proceed with
inserting into the DashMap if the reservation succeeds, and decrement/release
the counter/semaphore on removal or on insertion failure; alternatively
implement the check-and-insert via DashMap's Entry API combined with an
independent atomic counter so the global limit cannot be exceeded concurrently.
Ensure you update the corresponding cleanup paths to decrement/release the
reservation.
In `@model_gateway/src/server.rs`:
- Around line 616-631: The realtime_routes Router currently only applies
middleware::auth_middleware and thus bypasses the concurrency limiter/queue used
by protected_routes; update realtime_routes to apply the same
concurrency-limiting (and optional wasm) middleware as protected_routes (or
simply merge these routes into protected_routes) so session/token creation is
throttled—specifically, add the same route_layer(s) used by protected_routes to
realtime_routes (the Router named realtime_routes) in addition to
axum::middleware::from_fn_with_state(auth_config.clone(),
middleware::auth_middleware).
In `@protocols/src/realtime_response.rs`:
- Around line 111-128: ResponseStatusError currently only has r#type and code so
the API's "message" is dropped; update the ResponseStatusError struct by adding
a pub message: Option<String> field (keeping serde_with::skip_serializing_none
and the existing derives) so deserialization preserves the error message; ensure
the field name matches the JSON key ("message") and leave
RealtimeResponseStatusDetails unchanged aside from using the updated
ResponseStatusError.
🧹 Nitpick comments (5)
protocols/src/realtime_events.rs (2)
138-146: Consider preserving the original parse error for debugging.When the initial parse fails, the original error is discarded. This makes debugging difficult when the JSON is valid but doesn't match any known variant. Consider logging or preserving the original error.
♻️ Optional improvement
pub fn from_json(json: &str) -> Result<Self, serde_json::Error> { match serde_json::from_str::<ClientEvent>(json) { Ok(event) => Ok(event), - Err(_) => { + Err(e) => { + tracing::debug!("Unknown client event, preserving as Unknown: {}", e); let value: serde_json::Value = serde_json::from_str(json)?; Ok(ClientEvent::Unknown(value)) } } }
182-186: Fallback mapping for Unknown variant could be misleading.Mapping
ClientEvent::UnknowntoRealtimeClientEvent::SessionUpdateas a fallback is documented but could cause subtle bugs if callers forget to check forUnknownfirst. The comment "callers should check Unknown first" is helpful but easy to miss.Consider adding a distinct
Unknownvariant toRealtimeClientEventif possible, or documenting this more prominently in the method's doc comment.protocols/src/realtime_session.rs (2)
204-209: Consider documenting variant order for other untagged enums.Like
Voice, these#[serde(untagged)]enums (Tracing,RealtimeToolChoice,MaxOutputTokens,Truncation) rely on variant order for correct deserialization. While the patterns are generally safe (simple type vs object), adding brief comments similar to theVoiceenum would improve maintainability.Also applies to: 323-328, 342-347, 385-390
181-183: Nit: Missing blank line before section comment.Minor formatting inconsistency - other sections have a blank line before their separator comment.
♻️ Suggested fix
fn default_output_modalities() -> Vec<OutputModality> { vec![OutputModality::Audio] } + // ============================================================================ // Tracingmodel_gateway/src/routers/openai/realtime/rest.rs (1)
29-73: Consider extracting shared handler logic to avoid drift.The three handlers share the same select/auth/forward/proxy flow; a small internal helper could reduce duplication and future divergence.
796948e to
75fde82
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@model_gateway/src/routers/openai/realtime/rest.rs`:
- Around line 146-167: The forward_post function sets the worker Authorization
header before copying client headers, allowing a client-supplied "authorization"
to override it; update forward_post (or the header whitelist in
should_forward_request_header) so Authorization cannot be overridden by client
headers — either remove "authorization" from the whitelist in
should_forward_request_header, or move the .header("Authorization", auth) call
in forward_post to after the loop that copies headers (ensuring the worker auth
is set last and cannot be clobbered).
75fde82 to
1d1bc7d
Compare
There was a problem hiding this comment.
Actionable comments posted: 2
🤖 Fix all issues with AI agents
In `@model_gateway/Cargo.toml`:
- Around line 114-115: Update the workspace Tokio dependency to at least 1.42.1
to address the broadcast-channel unsoundness (RUSTSEC-2025-0023): modify the
workspace Cargo.toml Tokio version entry (the one referenced by
tokio-tungstenite and tokio-util) to >=1.42.1, run cargo update and cargo test;
additionally audit code that handles tungstenite::Message payloads (places using
tokio-tungstenite 0.26, functions/methods that match or construct Message
variants) and change handling from Vec<u8>/String to the new Bytes/Utf8Bytes
types (adjust pattern matches, conversions, and any usages of
Message::Text/Message::Binary accordingly).
In `@model_gateway/src/routers/openai/realtime/registry.rs`:
- Around line 211-287: The reaper in start_reaper currently evicts any entry
older than max_age and deducts counters using the stale_id list lengths, which
can remove active sessions and cause underflow if concurrent deletions occur;
update the logic to (1) when building stale_session_ids/stale_call_ids only
consider entries whose state is not the active/connected state (e.g., check
entry.state != SessionState::Connected or the equivalent enum variant used in
your code), and (2) instead of using
stale_session_ids.len()/stale_call_ids.len() to adjust session_count/call_count,
increase a local removed_sessions/removed_calls counter only when
registry.sessions.remove(id) / registry.calls.remove(id) returns Some and then
subtract that actual removed count (or clamp to current value) from
session_count/call_count to avoid underflow. Ensure you reference start_reaper,
registry.sessions, registry.calls, entry.state, entry.cancel_token,
session_count and call_count when making the changes.
|
Hi @pallasathena92, 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 |
1d1bc7d to
3e342ce
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/app_context.rs`:
- Line 305: RealtimeRegistry is instantiated with hardcoded defaults via
RealtimeRegistry::new(), so make its capacity configurable by constructing it
from the RouterConfig (or via a builder) instead of new(); replace the
Arc::new(RealtimeRegistry::new()) at realtime_registry with a call that consumes
limits from RouterConfig (e.g., RealtimeRegistry::from_config(...) or
RealtimeRegistry::with_limits(...)) or inject a prebuilt RealtimeRegistry from
the context builder so session/call caps come from RouterConfig and not a fixed
default.
3e342ce to
8be8a08
Compare
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 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/registry.rs`:
- Around line 254-287: The reaper currently snapshots
stale_session_ids/stale_call_ids and then removes by key, which can evict a
newly re-registered entry with the same id; instead, when iterating ids (or when
attempting removal) perform the staleness check inside the same critical section
before canceling/removing: for registry.sessions and registry.calls, replace the
two-phase snapshot/remove with a per-key remove-attempt that locks/queries the
entry and checks entry.state != ConnectionState::Connected &&
now.duration_since(entry.created_at) > max_age immediately before calling
registry.sessions.remove(id)/registry.calls.remove(id) and
entry.cancel_token.cancel(), so only entries that are still stale at removal
time are canceled and counted (update uses of stale_session_ids/stale_call_ids,
sessions_reaped, calls_reaped accordingly).
- Around line 53-55: The module relies on atomic counters and DashMap operations
to keep capacity counters consistent with map membership; add explicit
INVARIANT: comments at each reservation and removal site (e.g., the code paths
that increment/decrement the atomic capacity counters and the DashMap
insert/remove calls) stating the expected relationship between the counter and
map (for example, "INVARIANT: counter == number of entries in DASH_MAP for this
shard/tenant" and "INVARIANT: any increment must be paired with a subsequent
DashMap::insert on success; any decrement only after DashMap::remove
failed/rolled-back"). Place these INVARIANT: markers adjacent to the reservation
function(s) and removal/rollback function(s) so future reviewers can audit the
atomic counter ↔ DashMap consistency assumptions; continue to reserve SAFETY:
only for unsafe soundness explanations.
In `@model_gateway/src/routers/openai/realtime/rest.rs`:
- Around line 144-154: The forward_post function currently only sets
Authorization and JSON body but must also forward whitelisted trace headers for
propagation; update the signature of forward_post to accept a headers
map/Reference (e.g., &http::HeaderMap or &HeaderMap) and inside forward_post
copy the specific headers "x-request-id", "traceparent", "tracestate", and
"x-correlation-id" from that headers collection onto the reqwest::RequestBuilder
(using .header(...) only when the header exists), then update the call site that
invokes forward_post (where headers is available) to pass the headers argument
so the upstream receives the trace headers.
- Around line 173-189: In proxy_response, only Content-Type is forwarded which
loses useful upstream headers; update proxy_response(resp: reqwest::Response) to
copy a safe whitelist of upstream headers (e.g. "x-request-id", any
"x-ratelimit-*" matches, "retry-after") from resp.headers() into the outgoing
response headers while explicitly excluding hop-by-hop and sensitive headers, or
alternatively document that header propagation is intentionally disallowed;
locate and modify the headers construction around the match on
resp.bytes().await to build the response header list from the filtered
resp.headers() rather than the current single content-type entry.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (9)
Cargo.tomlmodel_gateway/Cargo.tomlmodel_gateway/src/app_context.rsmodel_gateway/src/routers/openai/mod.rsmodel_gateway/src/routers/openai/realtime/mod.rsmodel_gateway/src/routers/openai/realtime/registry.rsmodel_gateway/src/routers/openai/realtime/rest.rsmodel_gateway/src/server.rsmodel_gateway/src/service_discovery.rs
Signed-off-by: yifeliu <yifengliu9@gmail.com>
8be8a08 to
d9d2efa
Compare
openai api spec: https://developers.openai.com/api/reference/resources/realtime
Fix: #240 #241 #242
Description
Problem
The model gateway has no support for OpenAI's Realtime API. Clients that need ephemeral tokens for browser-safe WebSocket or WebRTC authentication (via /v1/realtime/client_secrets, /v1/realtime/sessions, or /v1/realtime/transcription_sessions) cannot route through the gateway, forcing them to connect directly to upstream workers and bypassing the gateway's load balancing, auth, and
circuit-breaker infrastructure.
Additionally, there is no in-memory tracking of active WebSocket sessions or WebRTC calls, which will be needed for upcoming WebSocket/WebRTC proxy phases to manage connection lifecycle, enforce capacity limits, and reap stale entries.
Solution
Introduce the Realtime API module (routers::openai::realtime) with two components:
interval.
The REST endpoints are wired into the axum router behind the existing auth and concurrency-limit middleware. The registry is created in AppContext and made available through AppState for use by future WebSocket/WebRTC handler phases.
Changes
Test Plan
Client secret: curl -X POST localhost:30000/v1/realtime/client_secrets -H "Authorization: Bearer $KEY" -d '{"session":{"type":"realtime","model":"gpt-4o-realtime-preview","audio":{"output":{"voice":"alloy"}}}}' → expect client_secret.value + expires_at
Session creation (legacy): curl -X POST http://localhost:30000/v1/realtime/sessions -H "Authorization: Bearer $KEY" -d '{"model":"gpt-4o-realtime-preview"}' → expect session object with client_secret.value
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Chores