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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
Member

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

CopilotAI 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
FileDescription
crates/aisix-provider-openai/src/wire.rsDefines OpenAI wire shapes and mapping logic for non-streaming + streaming chunks.
crates/aisix-provider-openai/src/lib.rsExposes OpenAiBridge and wires module structure.
crates/aisix-provider-openai/src/bridge.rsImplements the Bridge trait using reqwest + deadline handling + SSE decoding.
crates/aisix-provider-openai/Cargo.tomlAdds async/streaming dependencies and dev deps for wiremock tests.
Cargo.lockLocks 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 freeto 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

@moonming