Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading
, '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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
Expand Up@@ -79,6 +79,10 @@ jobs:
# every Config::load_from_path test as `redis_url` and breaks
# the entire `aisix-core::config::tests` module.
CACHE_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-ratelimit/tests/redis_integration.rs
# (shared cluster-level counters, #798). Same skip-if-unset /
# no-AISIX_-prefix rules as CACHE_TEST_REDIS_URL above.
RATELIMIT_TEST_REDIS_URL: redis://127.0.0.1:6379
# Picked up by crates/aisix-admin/tests/etcd_integration.rs.
# Same skip-if-unset pattern as the Redis case above.
ADMIN_TEST_ETCD_URL: http://127.0.0.1:2379
Expand Down
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

18 changes: 18 additions & 0 deletions config.example.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -91,6 +91,24 @@ cache:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel

# Rate-limit counter backend (api7/AISIX-Cloud#798).
#
# `memory` (default) keeps counters in this process, so a cluster of N
# replicas enforces N× every configured limit. `redis` shares the
# counters across every replica via one Redis, so the whole cluster
# enforces ONE global window — set this on multi-replica deployments.
# May point at the same Redis as `cache` (keys are namespaced
# `aisix:rl:`). On a Redis outage the limiter fails open to per-replica
# in-memory counting (logged) so traffic keeps flowing.
ratelimit:
backend: "memory" # memory | redis
# redis:
# url: "redis://127.0.0.1:6379"
# mode: "single" # single | cluster | sentinel
# Seconds before an unreleased concurrency slot (crashed replica /
# hung upstream) is reclaimed. Redis backend only.
# concurrency_ttl_secs: 300

# Models, API keys, provider keys, guardrails, cache policies, and
# observability exporters are NOT defined in this file. They are stored
# in etcd and managed via the Admin API (see docs/api-admin.md). This
Expand Down
6 changes: 6 additions & 0 deletions config.managed.yaml
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,12 @@
# `AISIX_CACHE__BACKEND`, etc. — every config field is reachable via
# `AISIX_<UPPER>__<UPPER>` (see crates/aisix-core/src/config.rs).
#
# For a multi-replica deployment, enable cluster-level rate limiting so
# the cluster enforces one global window instead of N× per replica
# (api7/AISIX-Cloud#798):
# AISIX_RATELIMIT__BACKEND=redis
# AISIX_RATELIMIT__REDIS__URL=redis://<host>:6379
#
# Subsequent boots re-use the mTLS bundle written under
# `managed.mtls_dir`.

Expand Down
144 changes: 144 additions & 0 deletions crates/aisix-core/src/config.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -43,6 +43,12 @@ pub struct Config {
pub observability: ObservabilityConfig,
#[serde(default)]
pub cache: CacheConfig,
/// Rate-limit counter backend. Defaults to per-process memory
/// (historical behaviour). Set `backend: redis` with a `redis` block
/// to share counters across every DP replica so a cluster enforces
/// one global window instead of one-per-replica (api7/AISIX-Cloud#798).
#[serde(default)]
pub ratelimit: RateLimitConfig,
/// Optional managed-mode configuration. When `managed.enabled = true`
/// the admin API and Playground endpoints are **not** bound — the DP
/// is a pure etcd reader driven by the aisix.cloud control plane.
Expand DownExpand Up@@ -559,6 +565,43 @@ impl RedisCacheConfig {
}
}

/// Rate-limit counter backend (api7/AISIX-Cloud#798).
///
/// `Memory` is the default: per-process fixed-window counters, so an
/// N-replica cluster enforces N× the configured limit. `Redis` shares
/// the counters across replicas via a single Redis so the whole cluster
/// enforces one global window. The `redis` block is required iff
/// `backend = redis` (validated at boot). Reuses [`RedisCacheConfig`]
/// for the connection shape; may point at the same Redis as `cache`
/// (keys are namespaced `aisix:rl:`).
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct RateLimitConfig {
pub backend: RateLimitBackend,
pub redis: Option<RedisCacheConfig>,
/// Seconds after which an unreleased concurrency slot is reclaimed
/// (crashed replica / hung upstream). Generous enough for a long
/// streaming response. Redis backend only.
pub concurrency_ttl_secs: u64,
}

impl Default for RateLimitConfig {
fn default() -> Self {
Self {
backend: RateLimitBackend::Memory,
redis: None,
concurrency_ttl_secs: 300,
}
}
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum RateLimitBackend {
Memory,
Redis,
}

impl Config {
/// Load + merge + validate.
///
Expand DownExpand Up@@ -658,6 +701,20 @@ impl Config {
"observability.metrics.prometheus.addr invalid socket address: {metrics_addr}"
)));
}
if self.ratelimit.backend == RateLimitBackend::Redis {
if self.ratelimit.redis.is_none() {
return Err(BootstrapError::Config(
"ratelimit.backend = redis requires a ratelimit.redis block".into(),
));
}
// A zero concurrency TTL would prune a slot in the same second
// it was taken, silently disabling concurrency limiting.
if self.ratelimit.concurrency_ttl_secs == 0 {
return Err(BootstrapError::Config(
"ratelimit.concurrency_ttl_secs must be > 0 for the redis backend".into(),
));
}
}
Ok(())
}
}
Expand DownExpand Up@@ -784,6 +841,93 @@ admin:
assert!(err.to_string().contains("admin.admin_keys"));
}

#[test]
fn ratelimit_defaults_to_memory_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Memory);
assert!(cfg.ratelimit.redis.is_none());
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 300);
}

#[test]
fn ratelimit_redis_backend_requires_redis_block() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("ratelimit.redis"));
}

#[test]
fn rejects_zero_concurrency_ttl_for_redis_backend() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 0
"#,
);
let err = Config::load_from_path(Some(f.path())).unwrap_err();
assert!(err.to_string().contains("concurrency_ttl_secs"));
}

#[test]
fn loads_ratelimit_redis_config() {
let f = write_yaml(
r#"
etcd:
endpoints: ["http://localhost:2379"]
proxy:
addr: "0.0.0.0:3000"
admin:
addr: "127.0.0.1:3001"
admin_keys: ["k1"]
ratelimit:
backend: "redis"
redis:
url: "redis://127.0.0.1:6379"
concurrency_ttl_secs: 120
"#,
);
let cfg = Config::load_from_path(Some(f.path())).unwrap();
assert_eq!(cfg.ratelimit.backend, RateLimitBackend::Redis);
assert_eq!(
cfg.ratelimit.redis.as_ref().unwrap().url,
"redis://127.0.0.1:6379"
);
assert_eq!(cfg.ratelimit.concurrency_ttl_secs, 120);
}

#[test]
fn rejects_invalid_bind_addr() {
let f = write_yaml(
Expand Down
3 changes: 2 additions & 1 deletion crates/aisix-core/src/lib.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -24,7 +24,8 @@ pub mod snapshot;

pub use config::{
AdminConfig, CacheBackend, CacheConfig, Config, EtcdConfig, EtcdTlsConfig, ManagedConfig,
ObservabilityConfig, ProxyConfig, RealIpConfig, TlsConfig,
ObservabilityConfig, ProxyConfig, RateLimitBackend, RateLimitConfig, RealIpConfig,
RedisCacheConfig, TlsConfig,
};
pub use error::{
AdminError, AdminErrorEnvelope, BootstrapError, ProxyError, ProxyErrorEnvelope, RateLimitScope,
Expand Down
23 changes: 12 additions & 11 deletions crates/aisix-proxy/src/chat.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -245,7 +245,7 @@ pub async fn chat_completions(
// current window state. We peek *after* the commit so
// remaining-requests reflects the post-dispatch tally.
let rl_limits = auth.key().rate_limit.clone().unwrap_or_default();
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits) {
if let Some(rl_status) = state.limiter.peek(&api_key_id, &rl_limits).await {
crate::render::inject_ratelimit_headers(&mut success.response, &rl_status);
state.metrics.set_rate_limit_remaining(
&api_key_id,
Expand DownExpand Up@@ -843,8 +843,9 @@ async fn dispatch(
&virtual_entry.id,
&virtual_entry.value,
);
let reservation =
crate::quota::enforce_rate_limit(state, auth, Some(&model_rl)).map_err(&with_model)?;
let reservation = crate::quota::enforce_rate_limit(state, auth, Some(&model_rl))
.await
.map_err(&with_model)?;

let now = created_ts();

Expand DownExpand Up@@ -1073,7 +1074,7 @@ async fn dispatch(
// permit was released here, letting a key capped at N run far more
// than N simultaneous streams (#450).
let post_stream_keys = reservation.keys();
let stream_concurrency_hold = reservation.into_stream_hold(Arc::clone(&state.limiter));
let stream_concurrency_hold = reservation.into_stream_hold();
// Capture everything the stream-completion callback needs so
// it can fire `emit_usage_event` once the terminal SSE chunk
// has yielded its `usage` block. Telemetry emission has to
Expand DownExpand Up@@ -1379,7 +1380,7 @@ async fn dispatch(
if let (Some(cache), Some(key)) = (policy_cache.as_ref(), cache_key.as_ref()) {
match cache.get(key).await {
Ok(Some(cached)) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
// #448: a cache hit is client-visible output just like a
// fresh upstream response, so it must run output guardrails
// before being returned — not bypass them.
Expand DownExpand Up@@ -1690,7 +1691,7 @@ async fn dispatch(
let provider_request_id = upstream.id.clone();
let provider_model_version = upstream.model.clone();
let finish_reason = finish_reason_label(&upstream.finish_reason);
reservation.commit_tokens(total);
reservation.commit_tokens(total).await;

// cp-api recomputes cost server-side from its pricing catalog when
// ingesting telemetry; the DP just records 0.0 on the wire.
Expand DownExpand Up@@ -1860,14 +1861,14 @@ async fn dispatch(
/// response. `reservation` is the SINGLE entry-level reservation taken
/// in `dispatch` — ensemble does not add per-sub-call reservations.
#[allow(clippy::too_many_arguments)]
async fn dispatch_ensemble<'a>(
state: &'a ProxyState,
async fn dispatch_ensemble(
state: &ProxyState,
snapshot: &aisix_core::AisixSnapshot,
virtual_entry: &aisix_core::ResourceEntry<aisix_core::Model>,
req: &ChatFormat,
request_id: &str,
created_ts: i64,
reservation: aisix_ratelimit::MultiReservation<'a, aisix_ratelimit::SystemClock>,
reservation: aisix_ratelimit::MultiReservation,
resolved_chain: &Arc<dyn aisix_guardrails::Guardrail>,
applied_guardrails: &[AppliedGuardrail],
mut bypass_reason: Option<String>,
Expand DownExpand Up@@ -2006,7 +2007,7 @@ async fn dispatch_ensemble<'a>(
),
};
let survivor_total: u64 = panel.iter().map(|p| u64::from(p.usage.total_tokens)).sum();
reservation.commit_tokens(survivor_total);
reservation.commit_tokens(survivor_total).await;
for (index, member) in panel.iter().enumerate() {
emit_panel_member(
member, index, /* blocked */ false, /* bypass */ "",
Expand All@@ -2030,7 +2031,7 @@ async fn dispatch_ensemble<'a>(
.sum();
let judge_usage = outcome.response.usage.clone();
let total_tokens = panel_total + u64::from(judge_usage.total_tokens);
reservation.commit_tokens(total_tokens);
reservation.commit_tokens(total_tokens).await;

// Emit one usage event per sub-call (each panel member + the judge),
// all sharing `request_id`. `attempt_kind` is `"panel"` / `"judge"`;
Expand Down
8 changes: 5 additions & 3 deletions crates/aisix-proxy/src/embeddings.rs
Original file line numberDiff line numberDiff line change
Expand Up@@ -321,7 +321,9 @@ async fn dispatch(
// and finalise RPM. Embeddings do report prompt_tokens via
// EmbeddingResponse.usage; thread it through so TPM works
// here even though other handlers commit 0.
reservation.commit_tokens(embed_resp.usage.total_tokens as u64);
reservation
.commit_tokens(embed_resp.usage.total_tokens as u64)
.await;
let provider_label = provider.to_ascii_lowercase();
// Capture the prompt_tokens count BEFORE moving the
// embed_resp into the JSON response — the handler needs
Expand All@@ -342,7 +344,7 @@ async fn dispatch(
// (`upstream_called: false` → handler skips emit per the
// chat.rs convention that we only attribute usage on a
// real upstream completion).
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
let env = ErrorEnvelope::new(msg, "not_implemented");
Ok(EmbedDispatchSuccess {
response: (StatusCode::NOT_IMPLEMENTED, Json(env)).into_response(),
Expand All@@ -357,7 +359,7 @@ async fn dispatch(
})
}
Err(e) => {
reservation.commit_tokens(0);
reservation.commit_tokens(0).await;
Err(ProxyError::Bridge(e))
}
}
Expand Down
Loading
Loading