feat(provider-openai): OpenAiBridge with streaming + non-streaming chat - #7
Merged
Conversation
First concrete Bridge implementation against the aisix-gateway trait. - wire.rs: OpenAI /chat/completions request and response wire types, plus the two mappers that round-trip between our ChatFormat / ChatResponse / ChatChunk and the upstream shape. Request extras flow through `#[serde(flatten)]` so seed/presence_penalty/etc. forward without the gateway having to know about them. - bridge.rs: OpenAiBridge owns a shared reqwest::Client. chat() does POST /chat/completions, parses the typed response, and applies the BridgeContext deadline via tokio::time::timeout. chat_stream() pipes bytes_stream() through the gateway's SseDecoder and yields ChatChunks via async_stream, terminating cleanly on the [DONE] sentinel. - Error mapping matches the BridgeError contract from PR #6: transport → Transport, non-2xx → UpstreamStatus (4xx passes through, 5xx collapses to 502 via http_status()), malformed JSON → UpstreamDecode, elapsed deadline → Timeout { elapsed_ms }. - `with_name()` lets OpenAI-compatible providers (DeepSeek today, Gemini-OAI later) reuse this transport with a distinct metrics label. 15 new unit tests across wire and bridge, 10 using wiremock: happy path (streaming + non-streaming), 429 pass-through, 500 pre-stream, malformed body → decode error, deadline → timeout, missing api_key → config error, SSE with role/content/finish_reason/[DONE], and resolve_base trailing- slash handling. 118 tests pass workspace-wide.
There was a problem hiding this comment.
Pull request overview
Adds the first concrete Bridge implementation (OpenAiBridge) for aisix-gateway, supporting both non-streaming and SSE-streaming OpenAI-compatible /chat/completions calls. This establishes a reusable transport layer intended to be shared by other OpenAI-compat providers (e.g., DeepSeek, Gemini-OAI).
Changes:
- Introduce OpenAI chat-completions wire/request/response types and mappers to/from
ChatFormat/ChatResponse/ChatChunk. - Implement
OpenAiBridgewith reqwest transport, error mapping, and SSE streaming viaSseDecoder. - Add provider crate dependencies and wiremock-backed unit tests.
Reviewed changes
Copilot reviewed 4 out of 5 changed files in this pull request and generated 7 comments.
Show a summary per file
| File | Description |
|---|---|
| crates/aisix-provider-openai/src/wire.rs | Defines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks. |
| crates/aisix-provider-openai/src/lib.rs | Exposes OpenAiBridge and wires module structure. |
| crates/aisix-provider-openai/src/bridge.rs | Implements the Bridge trait using reqwest + deadline handling + SSE decoding. |
| crates/aisix-provider-openai/Cargo.toml | Adds async/streaming dependencies and dev deps for wiremock tests. |
| Cargo.lock | Locks newly introduced dependencies (wiremock transitive deps, async-stream, etc.). |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
| tracing.workspace = true | ||
| tokio.workspace = true | ||
| futures.workspace = true | ||
| async-stream = "0.3" |
Comment on lines
+202
to
+203
| .map(|r| finish_reason(Some(r))) | ||
| .filter(|_| c.finish_reason.is_some()), |
Comment on lines
+80
to
+85
| fn resolve_base(model: &aisix_core::Model) -> String { | ||
| match model.base_url() { | ||
| Some(b) if !b.trim().is_empty() => b.trim_end_matches('/').to_string(), | ||
| _ => OPENAI_DEFAULT_BASE.to_string(), | ||
| } | ||
| } |
| if s.len() <= n { | ||
| s.to_string() | ||
| } else { | ||
| format!("{}…", &s[..n]) |
Comment on lines
+104
to
+108
| async fn map_http_error(status: StatusCode, resp: reqwest::Response) -> BridgeError { | ||
| let message = resp.text().await.unwrap_or_default(); | ||
| BridgeError::UpstreamStatus { | ||
| status: status.as_u16(), | ||
| message: truncate(&message, 1024), |
Comment on lines
+205
to
+227
| let resp = with_deadline(ctx.deadline, started, async move { | ||
| client | ||
| .post(&url) | ||
| .header(header::AUTHORIZATION, format!("Bearer {key}")) | ||
| .header(header::CONTENT_TYPE, "application/json") | ||
| .header(header::ACCEPT, "text/event-stream") | ||
| .header("x-aisix-request-id", &request_id) | ||
| .json(&body) | ||
| .send() | ||
| .await | ||
| .map_err(|e| BridgeError::Transport(e.to_string())) | ||
| }) | ||
| .await?; | ||
|
|
||
| let status = resp.status(); | ||
| if !status.is_success() { | ||
| return Err(map_http_error(status, resp).await); | ||
| } | ||
|
|
||
| let byte_stream = resp.bytes_stream(); | ||
| let stream = build_chunk_stream(byte_stream); | ||
| Ok(Box::pin(stream)) | ||
| } |
Comment on lines
+236
to
+257
| async_stream::try_stream! { | ||
| let mut decoder = SseDecoder::new(); | ||
| let mut stream = Box::pin(byte_stream); | ||
| while let Some(next) = stream.next().await { | ||
| let chunk = next.map_err(|e| BridgeError::Transport(e.to_string()))?; | ||
| for event in decoder.feed(chunk.as_ref()) { | ||
| match event { | ||
| SseEvent::Done => return, | ||
| SseEvent::Data(payload) => { | ||
| let parsed: OpenAiStreamChunk = serde_json::from_str(&payload) | ||
| .map_err(|e| BridgeError::UpstreamDecode(e.to_string()))?; | ||
| yield stream_chunk_into_chat_chunk(parsed); | ||
| } | ||
| } | ||
| } | ||
| } | ||
| if let Some(SseEvent::Data(payload)) = decoder.finish() { | ||
| let parsed: OpenAiStreamChunk = serde_json::from_str(&payload) | ||
| .map_err(|e| BridgeError::UpstreamDecode(e.to_string()))?; | ||
| yield stream_chunk_into_chat_chunk(parsed); | ||
| } | ||
| } |
moonming
added a commit
that referenced
this pull request
May 18, 2026
…ped reads, redact 5xx message, Vertex content-type guard Five concrete fixes from the Copilot inline review on PR #323. Two stale comments (#3, #4 — already fixed in commit 3) are skipped. **#1+#7 — Azure OpenAI-compatible code preservation.** Azure's envelope omits `error.type` and carries only `error.code`. The bridge previously put the upstream code into `view.kind` and left `view.code` as `None`. For OpenAI-compat tokens Azure inherits unchanged (e.g. `rate_limit_exceeded`), this meant downstream OpenAI clients received `error.type=rate_limit_exceeded` but `error.code=null` — exactly the SDK-retry break issue #322 is about. Fix: - Azure parser populates BOTH `view.kind` AND `view.code` from the upstream `error.code` field. - `render_openai_envelope`'s AzureOpenAI branch now prefers the translation-table-derived code (so explicit Azure tokens like `DeploymentNotFound` → `model_not_found` still win), falling back to `view.code` for OpenAI-compat pass-through. **#2 — Drain the response stream after hitting the cap.** `read_body_capped` previously broke out of the read loop the moment `limit` bytes were buffered. With reqwest/hyper that leaves unread bytes in the response and prevents connection reuse — during a burst of upstream errors the gateway would churn TCP connections instead of recycling the keep-alive pool. Fix: keep iterating the stream, discarding chunks past the cap. Memory stays bounded by `limit`. **#5 — Redact upstream `error.message` on 5xx.** The 5xx branch of `render_bridge_upstream_envelope` was forwarding `BridgeError::UpstreamStatus.message` verbatim — which for OpenAI / Anthropic comes from the parsed upstream `error.message`. Upstream 5xx bodies routinely embed operator-internal detail (engine names, shard ids, queue depth). Fix: on 5xx, emit a canned `"upstream returned {status}"` message; the full upstream body remains in operator logs via tracing. **#6 — Stale "follow-up" comment.** The docstring on `render_bridge_upstream_envelope` claimed cross-wire translation would ship in a follow-up, but it already shipped in commit 2. Rewrite the comment to describe current behaviour (4xx → `error_translate`; 5xx → canned envelope; `Unknown` wire → legacy generic envelope). **#8 — Content-type guard on Vertex (and Azure, while at it).** `capture_upstream_error_http` already gates serde parsing on `Content-Type: application/json` so a 64 KB HTML error page from a fronting WAF doesn't waste CPU on a doomed JSON parse. The Vertex and Azure bridges call serde directly because they need a custom parse path (canned message for redaction) — same guard now applies. Promoted `content_type_is_json` and added a `response_is_json` helper to the gateway's public surface; both bridges call it before `parse_*_error_*`. New tests: - `upstream_openai_5xx_with_json_envelope_collapses_and_redacts_message` pins the 5xx redaction (asserts `engine offline` / `shard 47` / `engine_overloaded` don't reach the customer envelope). - `chat_429_preserves_openai_compatible_code_for_sdk_retry` (Azure) pins that `parsed.code` carries the OpenAI-compat upstream code. - `chat_400_non_json_body_skips_envelope_parse` (Azure) and `chat_gemini_non_json_body_skips_envelope_parse` (Vertex) pin the new content-type guard. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
First concrete `Bridge` implementation against the `aisix-gateway`
trait. Also serves as the reusable transport for OpenAI-compatible
providers (DeepSeek today, Gemini-OAI later).
the mappers that round-trip `ChatFormat`/`ChatResponse`/`ChatChunk`
against the upstream shape. Request extras flow through
`#[serde(flatten)]` so `seed`, `presence_penalty`, etc. forward.
`chat()` does `POST /chat/completions` with a typed body +
tokio::time::timeout for deadlines. `chat_stream()` pipes
`bytes_stream()` through the gateway's `SseDecoder` and yields
`ChatChunk`s via `async_stream`, terminating on `[DONE]`.
non-2xx → UpstreamStatus, malformed JSON → UpstreamDecode, elapsed
deadline → Timeout { elapsed_ms }.
distinct metrics label.
Test plan