Skip to content

feat(provider-openai): OpenAiBridge with streaming + non-streaming chat - #7

Merged
moonming merged 1 commit into
mainfrom
feat/provider-openai
Apr 17, 2026
Merged

feat(provider-openai): OpenAiBridge with streaming + non-streaming chat#7
moonming merged 1 commit into
mainfrom
feat/provider-openai

Conversation

@moonming

Copy link
Copy Markdown
Collaborator

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).

  • wire.rs — OpenAI `/chat/completions` request/response types plus
    the mappers that round-trip `ChatFormat`/`ChatResponse`/`ChatChunk`
    against the upstream shape. Request extras flow through
    `#[serde(flatten)]` so `seed`, `presence_penalty`, etc. forward.
  • bridge.rs — `OpenAiBridge` owns a shared `reqwest::Client`.
    `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]`.
  • Error mapping follows the `BridgeError` contract: transport → Transport,
    non-2xx → UpstreamStatus, malformed JSON → UpstreamDecode, elapsed
    deadline → Timeout { elapsed_ms }.
  • `with_name()` lets OpenAI-compat providers reuse this transport with a
    distinct metrics label.

Test plan

  • 15 new unit tests (10 wiremock-backed)
    • non-streaming happy path
    • 429 pass-through
    • 500 before stream starts
    • malformed body → decode error
    • deadline → timeout
    • missing api_key → config error
    • SSE roundtrip (role/content/finish_reason/[DONE])
    • resolve_base trailing-slash handling
  • `cargo test --workspace` — 118 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

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.
Copilot AI review requested due to automatic review settings April 17, 2026 06:14
@moonming
moonming merged commit e253018 into main Apr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/provider-openai branch April 17, 2026 06:19

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 OpenAiBridge with reqwest transport, error mapping, and SSE streaming via SseDecoder.
  • 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>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants