feat(ratelimit): two-phase RPM/TPM/concurrency limiter + proxy wiring - #13
Conversation
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.
There was a problem hiding this comment.
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
| File | Description |
|---|---|
| crates/aisix-ratelimit/src/window.rs | Fixed-window counter with check/increment, add, and retry-after math + unit tests. |
| crates/aisix-ratelimit/src/limiter.rs | Two-phase per-key limiter and RAII reservation, including concurrency and token accounting + unit tests. |
| crates/aisix-ratelimit/src/lib.rs | Public crate surface + module wiring/exports. |
| crates/aisix-ratelimit/src/error.rs | Error taxonomy used by proxy to build 429 responses. |
| crates/aisix-ratelimit/src/clock.rs | Clock trait + SystemClock/TestClock to support deterministic unit tests. |
| crates/aisix-proxy/src/state.rs | Adds Arc<Limiter> to proxy state and a constructor variant intended for tests. |
| crates/aisix-proxy/src/chat.rs | Calls limiter pre-commit before upstream dispatch; commits tokens after completion. |
| crates/aisix-proxy/src/error.rs | Adds ProxyError::RateLimit mapping to 429 + optional Retry-After header. |
| crates/aisix-proxy/src/lib.rs | Extends 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.
| // 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); |
There was a problem hiding this comment.
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.
| // 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}"), | |
| }; |
| fn system_clock_returns_positive_now() { | ||
| assert!(SystemClock.unix_secs() > 0); |
There was a problem hiding this comment.
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.
| fn system_clock_returns_positive_now() { | |
| assert!(SystemClock.unix_secs() > 0); | |
| fn system_clock_unix_secs_does_not_panic() { | |
| let _ = SystemClock.unix_secs(); |
| // 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); | ||
| } | ||
| } |
There was a problem hiding this comment.
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.
| // 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); | |
| } | |
| } | |
| // 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. | |
| if let Some(max) = limits.rpd { | |
| if let WindowCheck::Full { retry_after_secs } = s.rpd.check_and_increment(now, 1, max) { | |
| return Err(RateLimitError::Requests { | |
| scope: RateLimitScope::Requests, | |
| retry_after_secs, | |
| }); | |
| } | |
| } | |
| if let Some(max) = limits.rpm { | |
| if let WindowCheck::Full { retry_after_secs } = s.rpm.check_and_increment(now, 1, max) { |
| // 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); |
There was a problem hiding this comment.
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).
| // 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); | |
| while let Some(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()), | |
| }; | |
| yield Ok::<_, Infallible>(ev); | |
| } | |
| // Emit the OpenAI-style [DONE] sentinel so clients that terminate | |
| // on it behave correctly. | |
| yield Ok::<_, Infallible>(Event::default().data("[DONE]")); | |
| reservation.commit_tokens(0); | |
| }; |
| /// 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, |
There was a problem hiding this comment.
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).
…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
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
`TestClock`.
`check_and_increment`/`add`/`is_exceeded`. Retry-after clamped to
≥1s.
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).
`retry_after_secs()` for the proxy to build `Retry-After`.
Proxy integration
tests).
Hub; on success commits `total_tokens` from the bridge response.
OpenAI-style `"type":"rate_limit_exceeded"`.
Test plan
math, concurrency drop safety, multi-key isolation, clock advance)