fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@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

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main) - #608

Merged
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api
Jun 15, 2026
Merged

fix: port ensemble streaming to the post-#607 MultiReservation API (unbreaks main)#608
moonming merged 1 commit into
mainfrom
fix/ensemble-stream-reservation-api

Conversation

@moonming

@moonmingmoonming commented Jun 15, 2026

Copy link
Copy Markdown
Member

🔴 Unbreaks main — build is currently red

main does not compile right now. #606 (ensemble streaming) was written against the pre-#607 reservation API and merged onto a main that already contained #607 ("cluster-level rate limiting via shared Redis"), which changed aisix_ratelimit::MultiReservation (dropped its generics/lifetime, made commit_tokens async, made into_stream_hold arg-less). There was no textual conflict (the two PRs touch different lines), so #606 merged CLEAN — but it's a semantic merge conflict: the merged tree doesn't build, and main's post-merge CI for #606 is failing. (My miss — I merged #606 without rebasing onto #607; filing this immediately.)

The fix

Port the streaming-ensemble reservation code in dispatch_ensemble to the new API (only the streaming arm; the single-upstream path, build_sse_stream, aisix-ratelimit, and aisix-core are untouched):

  • MultiReservation<'a, SystemClock>MultiReservation.
  • into_stream_hold(Arc::clone(&state.limiter))into_stream_hold().
  • Also fixes a latent billing bug: the error-exit helper bill_panel_and_fail was a sync closure calling the now-asynccommit_tokens with no .await — so post-feat(ratelimit): cluster-level rate limiting via shared Redis #607 the commit was a dropped future and the surviving panel tokens were never billed on streaming error exits. Split into a sync survivor_total + emit_panel_then_fail and .await reservation.commit_tokens(...) inline at each of the 6 error exits (mirroring the non-streaming ensemble path). One commit per request, every survivor billed exactly once — now actually executed.

Verification (local, against this branch = current main + the fix)

  • cargo check -p aisix-proxy — clean (main compiles again).
  • cargo clippy -p aisix-core -p aisix-proxy --all-targets -- -D warnings — zero warnings.
  • cargo test -p aisix-core — 263 pass; cargo test -p aisix-proxy — 440 pass (incl. ensemble_streams_judge_synthesis, ensemble_streaming_judge_connect_failure_still_bills_panel, and the rest of the 20 ensemble tests).
  • schema-drift gate (dump-schema + git diff --exit-code schemas/) — clean.

36/36-line change, isolated to chat.rs's streaming-ensemble arm. Please merge once CI is green to restore main.

Refs #601.

Summary by CodeRabbit

  • Refactor
    • Improved token billing and commitment process for streaming error scenarios.
    • Enhanced concurrency hold mechanism for streaming operations.

#606 (ensemble streaming) was written against the pre-#607 reservation
API and merged onto a main that already had #607's cluster-Redis change,
which dropped MultiReservation's generics/lifetime, made commit_tokens
async, and made into_stream_hold arg-less. The semantic merge conflict
broke main's build (no textual conflict, so it merged CLEAN).
Port the streaming-ensemble reservation code to the new API:
- MultiReservation<'a, SystemClock> -> MultiReservation.
- into_stream_hold(arg) -> into_stream_hold().
- The error-exit billing helper was a sync closure calling the now-async
commit_tokens with no .await, so the commit was a dropped future and
panel tokens were never billed on streaming error exits. Split into a
sync survivor_total + emit_panel_then_fail helper and .await
commit_tokens inline at each of the 6 exits (mirroring the
non-streaming path). One commit per request, every survivor billed
once -- now actually executed.
Verified: cargo check + clippy -D warnings clean, 263 core + 440 proxy
tests green (incl. both streaming e2e), schema-drift gate clean.
Refs #601.
@coderabbitai

coderabbitaiBot commented Jun 15, 2026

Copy link
Copy Markdown

Review Change Stack

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Free

Run ID: 84b1aa63-6566-46e0-aa01-f1a978c43a2f

📥 Commits

Reviewing files that changed from the base of the PR and between 9bd1e55 and ad0253d.

📒 Files selected for processing (1)
  • crates/aisix-proxy/src/chat.rs

📝 Walkthrough

Walkthrough

In dispatch_ensemble's streaming path, the synchronous bill_panel_and_fail closure is replaced with an async emit_panel_then_fail closure that awaits reservation.commit_tokens(survivor_total(&panel)) before emitting per-panel-member usage events and returning a DispatchFailure. All streaming error branches adopt this closure. The into_stream_hold call-site drops its explicit limiter argument.

Changes

Streaming Ensemble Async Token Commit

Layer / File(s)Summary
emit_panel_then_fail closure and all streaming error branches
crates/aisix-proxy/src/chat.rs
Defines an async emit_panel_then_fail closure that computes survivor_total(&panel), awaits reservation.commit_tokens(...), emits panel-member usage events, and returns DispatchFailure. Applies it to every streaming error exit: InsufficientPanel, judge-phase, missing judge model entry, unresolvable provider key, unregistered bridge, and failed chat_stream connect.
Concurrency hold simplification
crates/aisix-proxy/src/chat.rs
Changes reservation.into_stream_hold(Arc::clone(&state.limiter)) to reservation.into_stream_hold(), removing the explicit limiter argument.

Estimated code review effort

🎯 2 (Simple) | ⏱️ ~10 minutes


Note

🎁 Summarized by CodeRabbit Free

Your organization has reached its limit of developer seats under the Pro Plan. For new users, CodeRabbit will generate a high-level summary and a walkthrough for each pull request. For a comprehensive line-by-line review, please add seats to your subscription by visiting https://app.coderabbit.ai/login.If you believe this is a mistake and have available seats, please assign one to the pull request author through the subscription management page using the link above.

Comment @coderabbitai help to get the list of available commands and usage tips.

@moonming
moonming merged commit 3e8827e into mainJun 15, 2026
8 checks passed
@moonming
moonming deleted the fix/ensemble-stream-reservation-api branch June 15, 2026 04:55
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.

1 participant

@moonming