feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13

Merged
moonming merged 1 commit into
mainfrom
feat/ratelimit
Apr 17, 2026
Merged

feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring#13
moonming merged 1 commit into
mainfrom
feat/ratelimit

Conversation

@moonming

Copy link
Copy Markdown
Member

Summary

Spec §3 rate limiting. Two-phase: RPM/RPD checked-and-incremented up
front (burst requests fail fast), TPM/TPD checked up front + added
post-deduct (token cost is only known after upstream completes),
concurrency enforced via an in-flight counter guarded by the same per-
key mutex.

`aisix-ratelimit` layout

  • clock.rs — `Clock` trait + `SystemClock` + deterministic
    `TestClock`.
  • window.rs — `FixedWindowCounter` with
    `check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
    ≥1s.
  • limiter.rs — `Limiter` + RAII `Reservation`. `pre_commit`
    does the full check sequence; `commit_tokens` records TPM/TPD and
    releases the concurrency permit. Dropping without commit still
    releases the permit (no leaks on panicking upstream paths).
  • error.rs — `RateLimitError` exposes `scope()` and
    `retry_after_secs()` for the proxy to build `Retry-After`.

Proxy integration

  • `ProxyState` now carries `Arc` (with_limiter override for
    tests).
  • `chat_completions` calls `pre_commit` before dispatching to the
    Hub; on success commits `total_tokens` from the bridge response.
  • `ProxyError::RateLimit` maps to 429 with `Retry-After` header and
    OpenAI-style `"type":"rate_limit_exceeded"`.

Test plan

  • 18 unit tests in aisix-ratelimit (counter rollover, retry-after
    math, concurrency drop safety, multi-key isolation, clock advance)
  • 2 end-to-end proxy tests via axum + wiremock
    • `rpm=1` → first 200, second 429 with `Retry-After`
    • `tpm=1000` + overshoot usage → first 200, second 429
  • `cargo test --workspace` — 210 tests pass
  • `cargo clippy --all-targets -- -D warnings` clean
  • `cargo fmt --check` clean
  • CI green across all 6 jobs

aisix-ratelimit implements the two-phase limiter described in spec §3:
RPM/RPD checked-and-incremented at pre-commit, TPM/TPD checked at
pre-commit + added at post-deduct (we only know token usage after the
upstream response lands), concurrency enforced via an in-flight counter
guarded by the same per-key mutex.
aisix-ratelimit layout:
- clock.rs: Clock trait + SystemClock (production) + TestClock
(deterministic stepper for unit tests).
- window.rs: FixedWindowCounter — a single second-granularity bucket
with check_and_increment / add / is_exceeded helpers. Retry-after
hint clamped to >=1s.
- limiter.rs: Limiter<C> + Reservation RAII guard. pre_commit does the
full check sequence; commit_tokens records TPM/TPD and releases the
concurrency permit. Dropping without commit still releases the
permit so panicking upstream paths don't leak in_flight capacity.
- error.rs: RateLimitError with scope() + retry_after_secs() for the
proxy layer to produce Retry-After.
Proxy integration:
- ProxyState now carries Arc<Limiter>; ProxyState::with_limiter() lets
tests inject a shared instance.
- chat_completions handler calls pre_commit before dispatching to the
Hub; on success commits total_tokens from the bridge response.
Streaming currently commits 0 (streaming token counting is a later
PR).
- ProxyError grows a RateLimit variant that maps to 429 with a
Retry-After header and OpenAI-shaped "type":"rate_limit_exceeded".
Tests:
- 18 unit tests in aisix-ratelimit (counter rollover, retry-after math,
concurrency drop safety, multi-key isolation, clock advance).
- 2 end-to-end proxy tests via axum + wiremock:
- rpm=1 → first 200, second 429 with a Retry-After header
- tpm=1000 + overshoot usage → first 200, second 429
210 tests pass workspace-wide.
CopilotAI review requested due to automatic review settings April 17, 2026 07:57
@moonming
moonming merged commit 929d57c into mainApr 17, 2026
9 checks passed
@moonming
moonming deleted the feat/ratelimit branch April 17, 2026 08:02

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 a new aisix-ratelimit crate implementing a two-phase (pre-commit + post-deduct) rate limiter for RPM/RPD and TPM/TPD, plus per-key concurrency limiting, and wires it into aisix-proxy so /v1/chat/completions enforces limits and returns OpenAI-style 429s with Retry-After.

Changes:

  • Introduce aisix-ratelimit (clock abstraction, fixed-window counters, limiter + reservation RAII, error type).
  • Integrate limiter into proxy state + chat handler; map limiter errors to 429 with optional Retry-After.
  • Add proxy E2E tests covering RPM and TPM limiting behavior.

Reviewed changes

Copilot reviewed 9 out of 9 changed files in this pull request and generated 5 comments.

Show a summary per file
FileDescription
crates/aisix-ratelimit/src/window.rsFixed-window counter with check/increment, add, and retry-after math + unit tests.
crates/aisix-ratelimit/src/limiter.rsTwo-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests.
crates/aisix-ratelimit/src/lib.rsPublic crate surface + module wiring/exports.
crates/aisix-ratelimit/src/error.rsError taxonomy used by proxy to build 429 responses.
crates/aisix-ratelimit/src/clock.rsClock trait + SystemClock/TestClock to support deterministic unit tests.
crates/aisix-proxy/src/state.rsAdds Arc<Limiter> to proxy state and a constructor variant intended for tests.
crates/aisix-proxy/src/chat.rsCalls limiter pre-commit before upstream dispatch; commits tokens after completion.
crates/aisix-proxy/src/error.rsAdds ProxyError::RateLimit mapping to 429 + optional Retry-After header.
crates/aisix-proxy/src/lib.rsExtends test helpers for rate_limit config; adds two E2E tests for rpm/tpm limiting.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment on lines +456 to +458
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()), make_req()).await;
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

These proxy rate-limit E2E tests use ProxyState::new() which constructs a Limiter backed by SystemClock, so whether the 2nd request is still in the same RPM/TPM window depends on wall-clock timing. This can become flaky around minute boundaries. Consider injecting a deterministic clock-backed limiter for tests (or otherwise freezing time) so the retry logic is stable.

Suggested change
// Second request within the same minute trips rpm=1 → 429.
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
// With a wall-clock-backed limiter, a minute boundary can fall between
// the first and second request. If that happens, the second request may
// be accepted into the next minute's bucket and the third request must
// then be rate-limited.
let resp = run(build_router(state.clone()),make_req()).await;
let resp = match resp.status(){
StatusCode::TOO_MANY_REQUESTS => resp,
StatusCode::OK => {
let resp = run(build_router(state.clone()),make_req()).await;
assert_eq!(resp.status(),StatusCode::TOO_MANY_REQUESTS);
resp
}
status => panic!("expected second request to be 200 or 429, got {status}"),
};

Copilot uses AI. Check for mistakes.
Comment on lines +62 to +63
fn system_clock_returns_positive_now() {
assert!(SystemClock.unix_secs() > 0);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

system_clock_returns_positive_now can be flaky in environments where the system clock is misconfigured (e.g., set before UNIX_EPOCH), because SystemClock::unix_secs() explicitly returns 0 on duration_since errors. Consider asserting only that the call doesn’t panic (or that it returns >= 0), or remove this test since it’s validating external environment state rather than your logic.

Suggested change
fnsystem_clock_returns_positive_now(){
assert!(SystemClock.unix_secs() > 0);
fnsystem_clock_unix_secs_does_not_panic(){
let _ = SystemClock.unix_secs();

Copilot uses AI. Check for mistakes.
Comment on lines +126 to +147
// Request limits — checked AND incremented.
if let Some(max) = limits.rpm {
if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) {
return Err(RateLimitError::Requests {
scope: RateLimitScope::Requests,
retry_after_secs,
});
}
}
if let Some(max) = limits.rpd {
if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) {
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0 {
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
if let Some(max) = limits.rpm {
let _ = s.rpm.check_and_increment(now, 0, max);
}
}

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

The RPM rollback logic in the RPD-limit failure path resets s.rpm to a brand-new counter, which discards the current window start and any prior RPM usage in this minute. That will undercount RPM after an RPD rejection and can allow more RPM than configured once RPD stops rejecting. Consider avoiding partial mutation (e.g., check RPD before incrementing RPM), or implement an explicit decrement/restore of the previous RPM count within the same window.

Suggested change
// Request limits — checked AND incremented.
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
// Compensate: we already incremented RPM above. Decrement
// it so the caller's retry on a different day still
// counts correctly. RPM would have rolled by then, so
// this is primarily defensive.
if s.rpm.current(now) > 0{
// Roll back the increment we just made.
s.rpm = FixedWindowCounter::new(MINUTE_SECS);
ifletSome(max) = limits.rpm{
let _ = s.rpm.check_and_increment(now,0, max);
}
}
// Request limits — checked AND incremented. Check the broader
// daily window before mutating the minute window so we never
// need to roll back RPM state on an RPD rejection.
ifletSome(max) = limits.rpd{
ifletWindowCheck::Full{ retry_after_secs } = s.rpd.check_and_increment(now,1, max){
returnErr(RateLimitError::Requests{
scope:RateLimitScope::Requests,
retry_after_secs,
});
}
}
ifletSome(max) = limits.rpm{
ifletWindowCheck::Full{ retry_after_secs } = s.rpm.check_and_increment(now,1, max){

Copilot uses AI. Check for mistakes.
Comment on lines +76 to 83
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req, &ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

In the streaming branch, calling reservation.commit_tokens(0) consumes and drops the reservation before the SSE body is actually streamed to completion, which releases the concurrency permit immediately. That means the concurrency limit is not enforced for the lifetime of a streaming request. To enforce concurrency correctly, keep the Reservation alive until the stream ends (e.g., wrap the stream in a struct that owns the reservation and commits/drops it in Drop/on completion).

Suggested change
// Streaming: we can't measure tokens before the stream ends, so
// commit zero up front to keep the reservation's drop-guard from
// silently counting nothing. A later PR will tally tokens as the
// stream runs; for now release the permit when the handler returns.
let upstream = bridge.chat_stream(&req,&ctx).await?;
reservation.commit_tokens(0);
let model_name = req.model.clone();
let sse_stream = build_sse_stream(upstream, model_name, now);
// Streaming: keep the reservation alive for the full lifetime of the
// SSE body so the concurrency permit is not released when the handler
// returns. We still commit zero tokens for now, but only once the
// stream finishes; if the client disconnects early, dropping the
// stream will drop the reservation and release the permit then.
let upstream = bridge.chat_stream(&req,&ctx).await?;
let _model_name = req.model.clone();
let sse_stream = async_stream::stream! {
let reservation = reservation;
futures::pin_mut!(upstream);
whileletSome(item) = upstream.next().await{
let ev = match item {
Ok(chunk) => {
let rendered = render_chunk(now, chunk);
match serde_json::to_string(&rendered){
Ok(json) => Event::default().data(json),
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
}
}
Err(err) => Event::default()
.event("error")
.data(err.to_string()),
};
yieldOk::<_,Infallible>(ev);
}
// Emit the OpenAI-style [DONE] sentinel so clients that terminate
// on it behave correctly.
yieldOk::<_,Infallible>(Event::default().data("[DONE]"));
reservation.commit_tokens(0);
};

Copilot uses AI. Check for mistakes.
Comment on lines +37 to +43
/// Alternative constructor for callers that want to share a preexisting
/// limiter (e.g. tests with a deterministic clock).
pub fn with_limiter(
snapshot: SnapshotHandle<AisixSnapshot>,
hub: Arc<Hub>,
limiter: Arc<Limiter>,
cfg: &ProxyConfig,

CopilotAIApr 17, 2026

Copy link

Choose a reason for hiding this comment

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

ProxyState::with_limiter is documented as enabling tests with a deterministic clock, but ProxyState stores Arc<Limiter> (defaults to Limiter<SystemClock>), so callers cannot pass Limiter<TestClock> here. Either adjust the docs, or change the field/type to allow injecting a non-SystemClock limiter (e.g., make Limiter non-generic by storing an Arc<dyn Clock>, or store a trait object for the limiter).

Copilot uses AI. Check for mistakes.
jarvis9443 added a commit that referenced this pull request Jun 2, 2026
…eaming)
Thread the resolved guardrail chain (as Arc) through the /v1/messages
dispatch paths and run output guardrails on the response:
- Non-streaming: cross-provider checks the bridge ChatResponse;
passthrough extracts response text (content blocks + raw content array
for tool_use) into a synthetic ChatResponse.
- Streaming: both the cross-provider SSE encoder path and the verbatim
Anthropic byte-passthrough accumulate assistant text and run the
guardrail at end-of-stream. Bytes are forwarded live (matching
/v1/chat/completions and LiteLLM's streaming guardrail), so a block is
signalled with a terminal Anthropic `error` (content_filter) event.
Completes the output side of #448#22; with this and the earlier input +
budget work, /v1/messages no longer bypasses the guardrail/quota
pipeline. The remaining findings (#6 count_tokens, #2/#13
reasoning_content, #24 guardrail-vs-rate-limit ordering) are accepted as
standard behavior (LiteLLM has the same gap).
Fixes#448
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