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
2 changes: 1 addition & 1 deletion Cargo.lock

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

2 changes: 1 addition & 1 deletion plugins/rust/python-package/rate_limiter/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "rate_limiter"
version = "0.1.10"
version = "0.1.11"
edition.workspace = true
authors.workspace = true
license.workspace = true
Expand Down
24 changes: 16 additions & 8 deletions plugins/rust/python-package/rate_limiter/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

> Author: ContextForge Contributors

Enforces rate limits per user, tenant, and tool across `tool_pre_invoke` and `prompt_pre_fetch` hooks. Supports pluggable counting algorithms (fixed window, sliding window, token bucket), an in-process memory backend (single-instance), and a Redis backend (shared across all gateway instances).
Enforces global per-user, per-user/per-server, tenant, and tool rate limits across `tool_pre_invoke` and `prompt_pre_fetch` hooks. Supports pluggable counting algorithms (fixed window, sliding window, token bucket), an in-process memory backend (single-instance), and a Redis backend (shared across all gateway instances).

## Runtime Requirements

Expand All @@ -12,8 +12,8 @@ This plugin depends on `cpex>=0.1.0,<0.2` and imports hook models from `cpex.fra

| Hook | When it runs |
|---|---|
| `tool_pre_invoke` | Before every tool call — checks `by_user`, `by_tenant`, `by_tool` |
| `prompt_pre_fetch` | Before every prompt fetch — checks `by_user`, `by_tenant`, `by_tool` |
| `tool_pre_invoke` | Before every tool call — checks `by_user`, `by_user_per_server`, `by_tenant`, `by_tool` |
| `prompt_pre_fetch` | Before every prompt fetch — checks `by_user`, `by_user_per_server`, `by_tenant`, `by_tool` |

If any configured dimension is exceeded, the plugin returns a violation with HTTP 429. All requests include `X-RateLimit-*` headers. The most restrictive active dimension is surfaced (e.g. if both user and tenant limits are active, the one closest to exhaustion is reported).

Expand All @@ -27,7 +27,8 @@ If any configured dimension is exceeded, the plugin returns a violation with HTT
- tool_pre_invoke
mode: enforce # enforce | permissive | disabled
config:
by_user: "30/m" # per-user limit across all tools
by_user: "300/m" # per-user limit across all servers
by_user_per_server: "60/m" # per-user limit for each MCP server
by_tenant: "300/m" # shared limit across all users in a tenant
by_tool: # per-tool overrides (applied on top of by_user)
search: "10/m"
Expand Down Expand Up @@ -55,6 +56,7 @@ If any configured dimension is exceeded, the plugin returns a violation with HTT
| Field | Type | Default | Description |
|---|---|---|---|
| `by_user` | string | `null` | Per-user rate limit, e.g. `"60/m"` |
| `by_user_per_server` | string | `null` | Per-user rate limit independently enforced for each MCP server, e.g. `"60/m"` |
| `by_tenant` | string | `null` | Per-tenant rate limit, e.g. `"600/m"` |
| `by_tool` | dict | `{}` | Per-tool overrides, e.g. `{"search": "10/m"}` |
| `algorithm` | string | `"fixed_window"` | Counting algorithm: `"fixed_window"`, `"sliding_window"`, or `"token_bucket"` |
Expand All @@ -75,6 +77,8 @@ If any configured dimension is exceeded, the plugin returns a violation with HTT

**Omitting a dimension** (e.g. no `by_tenant`) means that dimension is unlimited — no counter is tracked for it.

`by_user` and `by_user_per_server` are independent and may be enabled together. For example, `by_user: "300/m"` plus `by_user_per_server: "60/m"` gives each user at most 60 requests per server and 300 requests across all servers. This prevents one user from multiplying their effective quota without relying on the shared `by_tenant` ceiling. When `by_user_per_server` is configured, requests without a safe non-empty `server_id` are blocked with `RATE_LIMIT_CONTEXT_MISSING` (HTTP 503) rather than bypassing the limit.

## Response headers

Every request (allowed or blocked) includes:
Expand Down Expand Up @@ -106,7 +110,7 @@ Stores a timestamp for every request in the current window. At each check, expir

### Token bucket

Each identity (user, tenant, tool) has a bucket that holds up to `count` tokens. Tokens refill at a steady rate of `count/window`. A request consumes one token. Bursts up to the bucket capacity are allowed; sustained rate above `count/window` is rejected. Useful for APIs where short spikes are acceptable but sustained overload is not.
Each identity (user, user/server, tenant, tool) has a bucket that holds up to `count` tokens. Tokens refill at a steady rate of `count/window`. A request consumes one token. Bursts up to the bucket capacity are allowed; sustained rate above `count/window` is rejected. Useful for APIs where short spikes are acceptable but sustained overload is not.

**Redis support:** `token_bucket` with `backend: redis` is fully supported. The plugin stores `{tokens, last_refill}` in a Redis hash per key and uses an atomic Lua script to refill and consume tokens in a single round-trip — the same pattern as the other two algorithms. This means `token_bucket` enforces a true cluster-wide limit in multi-instance deployments.

Expand Down Expand Up @@ -178,11 +182,12 @@ When the plugin context carries a `tenant_id`, every dimension key is prefixed w

```
rl:{tenant_id}:user:{email}:{window_seconds}
rl:{tenant_id}:user:{email}:server:{server_id}:{window_seconds}
rl:{tenant_id}:tenant:{tenant_id}:{window_seconds}
rl:{tenant_id}:tool:{tool_name}:{window_seconds}
```

When `tenant_id` is absent (single-tenant deployments), the prefix is omitted and keys revert to the pre-tenant-scoping layout (`rl:user:{email}:{window}`), so single-tenant behaviour is unchanged.
When `tenant_id` is absent (single-tenant deployments), the tenant prefix is omitted. Global user keys retain the existing layout (`rl:user:{email}:{window}`); per-server keys use `rl:user:{email}:server:{server_id}:{window}`.

**Upgrade note:** the first deploy of the tenant-scoping change causes counters under `rl:user:*` / `rl:tool:*` to be orphaned while new writes land at `rl:{tenant}:user:*`. Counters effectively reset once for all in-flight windows — non-event for typical second/minute windows.

Expand All @@ -192,7 +197,8 @@ When `tenant_id` is absent (single-tenant deployments), the prefix is omitted an

```yaml
config:
by_user: "60/m"
by_user: "300/m"
by_user_per_server: "60/m"
by_tenant: "600/m"
```

Expand Down Expand Up @@ -267,6 +273,8 @@ result.metadata["rate_limiter"] = {

## Migration Note

Version `0.1.11` adds the optional `by_user_per_server` dimension. Existing configurations and global `by_user` Redis keys remain unchanged. Enabling the new field creates fresh user/server counters; deploy `0.1.11` to every replica before enabling it so all replicas enforce the same dimensions. Use the stable catalog UUID carried in `GlobalContext.server_id`; display names or slugs can reset quotas when renamed.

Version `0.1.7` is a **breaking change** for any existing consumer reading rate-limit metadata:

- The old flat, unconditional `result.metadata` write (the engine's `meta` dict — `limited`, `remaining`, `reset_in`, `dimensions` — written on every allowed/not-limited call regardless of trace context) is now gated on a valid `trace_id` and replaced by the namespaced `rate_limiter` key containing only `allowed`/`throttled`/`backend` (the engine's own fields are not folded in — see "Returned Metadata" above for why).
Expand All @@ -291,5 +299,5 @@ Without `shutdown`, the cached Redis connection would leak across plugin re-inst
| Fixed window allows up to 2× limit at window boundary | LOW | Use `sliding_window` algorithm, or use `by_user` with headroom |
| `by_tool` matching is case-sensitive | LOW | Fixed — tool names are normalised with `.strip().lower()` |
| Whitespace-only user identity bypasses anonymous bucket | LOW | Fixed — `_extract_user_identity` strips whitespace and falls back to `'anonymous'` |
| No per-server limits (`server_id` dimension missing) | LOW | Not implemented |
| Per-server user limits require a stable server ID | MEDIUM | Use the catalog UUID propagated as `GlobalContext.server_id` |
| No config hot-reload — rate string changes require restart | LOW | Not implemented |
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
description: "Rate limiting by user/tenant/tool — memory (single-process) or Redis (shared across instances)"
description: "Rate limiting by user, user/server, tenant, and tool — memory or Redis backends"
author: "ContextForge Contributors"
version: "0.1.10"
version: "0.1.11"
kind: "cpex_rate_limiter.rate_limiter.RateLimiterPlugin"
available_hooks:
- "prompt_pre_fetch"
- "tool_pre_invoke"
default_configs:
by_user: "60/m"
by_user_per_server: null
by_tenant: "600/m"
by_tool: {}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ def _parse_rate(rate: str) -> tuple[int, int]:
class RateLimiterConfig:
__slots__ = (
"by_user",
"by_user_per_server",
"by_tenant",
"by_tool",
"algorithm",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,7 @@ class RateLimiterEngine:
`redis_ur` instead of `redis_url`) surface visibly instead of being
silently ignored.
"""
def check(self, user: builtins.str, tenant: typing.Optional[builtins.str], tool: builtins.str, now_unix: builtins.int, include_retry_after: builtins.bool, context_prefix: typing.Optional[builtins.str]) -> tuple[builtins.bool, dict, dict]:
def check(self, user: builtins.str, tenant: typing.Optional[builtins.str], tool: builtins.str, now_unix: builtins.int, include_retry_after: builtins.bool, context_prefix: typing.Optional[builtins.str], server_id: typing.Optional[builtins.str] = None) -> tuple[builtins.bool, dict, dict]:
r"""
High-level check: builds dimension keys internally, evaluates, and
returns pre-built Python dicts for headers and metadata.
Expand All @@ -135,7 +135,7 @@ class RateLimiterEngine:
path is intended for the memory backend. The `debug_assert` below
guards against accidental misuse.
"""
def check_async(self, user: builtins.str, tenant: typing.Optional[builtins.str], tool: builtins.str, now_unix: builtins.int, include_retry_after: builtins.bool, context_prefix: typing.Optional[builtins.str]) -> typing.Any:
def check_async(self, user: builtins.str, tenant: typing.Optional[builtins.str], tool: builtins.str, now_unix: builtins.int, include_retry_after: builtins.bool, context_prefix: typing.Optional[builtins.str], server_id: typing.Optional[builtins.str] = None) -> typing.Any:
r"""
Async variant of `check()` for Redis-backed deployments.

Expand Down
33 changes: 29 additions & 4 deletions plugins/rust/python-package/rate_limiter/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ impl Algorithm {
#[derive(Debug, Clone)]
pub struct EngineConfig {
pub by_user: Option<RateLimit>,
pub by_user_per_server: Option<RateLimit>,
pub by_tenant: Option<RateLimit>,
/// Normalised key → limit. Keys are already `.trim().to_lowercase()`.
pub by_tool: HashMap<String, RateLimit>,
Expand All @@ -119,6 +120,7 @@ impl EngineConfig {
/// that are relevant to the Rust engine — strict subset per IFACE-04).
pub fn new(
by_user: Option<&str>,
by_user_per_server: Option<&str>,
by_tenant: Option<&str>,
by_tool: HashMap<String, String>,
algorithm: &str,
Expand All @@ -131,6 +133,14 @@ impl EngineConfig {
})
})
.transpose()?;
let by_user_per_server = by_user_per_server
.map(|rate| {
parse_rate(rate).map_err(|err| ConfigError::FieldError {
field: format!("by_user_per_server={rate:?}"),
message: err.to_string(),
})
})
.transpose()?;
let by_tenant = by_tenant
.map(|rate| {
parse_rate(rate).map_err(|err| ConfigError::FieldError {
Expand All @@ -155,6 +165,7 @@ impl EngineConfig {
.ok_or_else(|| ConfigError::InvalidAlgorithm(algorithm.to_string()))?;
Ok(Self {
by_user,
by_user_per_server,
by_tenant,
by_tool,
algorithm,
Expand Down Expand Up @@ -256,9 +267,17 @@ mod tests {
by_tool.insert("Search".to_string(), "10/m".to_string());
by_tool.insert(" Summarise ".to_string(), "5/m".to_string());

let cfg = EngineConfig::new(Some("30/m"), Some("300/m"), by_tool, "fixed_window").unwrap();
let cfg = EngineConfig::new(
Some("30/m"),
Some("15/m"),
Some("300/m"),
by_tool,
"fixed_window",
)
.unwrap();

assert_eq!(cfg.by_user.unwrap().count, 30);
assert_eq!(cfg.by_user_per_server.unwrap().count, 15);
assert_eq!(cfg.by_tenant.unwrap().count, 300);
// Keys must be normalised
assert!(cfg.by_tool.contains_key("search"));
Expand All @@ -269,19 +288,25 @@ mod tests {

#[test]
fn engine_config_all_none_is_valid() {
let cfg = EngineConfig::new(None, None, HashMap::new(), "sliding_window").unwrap();
let cfg = EngineConfig::new(None, None, None, HashMap::new(), "sliding_window").unwrap();
assert!(cfg.by_user.is_none());
assert!(cfg.by_user_per_server.is_none());
assert!(cfg.by_tenant.is_none());
assert!(cfg.by_tool.is_empty());
}

#[test]
fn engine_config_invalid_rate_propagates_error() {
assert!(EngineConfig::new(Some("bad"), None, HashMap::new(), "fixed_window").is_err());
assert!(
EngineConfig::new(Some("bad"), None, None, HashMap::new(), "fixed_window").is_err()
);
assert!(
EngineConfig::new(None, Some("bad"), None, HashMap::new(), "fixed_window").is_err()
);
}

#[test]
fn engine_config_invalid_algorithm_propagates_error() {
assert!(EngineConfig::new(None, None, HashMap::new(), "leaky_bucket").is_err());
assert!(EngineConfig::new(None, None, None, HashMap::new(), "leaky_bucket").is_err());
}
}
Loading
Loading