mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-28 16:07:05 +00:00
feat(api): wire opt-in per-client S3 API rate limiting (429 + Retry-After) (#4895)
feat(api): wire opt-in per-client S3 API rate limiting (backlog#1191) RustFS shipped three rate-limiter implementations and none was wired to any request path: the tower layer never returned 429 (its over-limit branch passed requests through) and was never instantiated, the console env switches only logged, and the Swift token bucket was never called. Replace them with one working, default-off implementation: - Rewrite rustfs/src/server/rate_limit.rs as a sharded per-client-IP token-bucket limiter (32 mutex shards instead of one global RwLock write per request), bounded at 100k tracked IPs with lossless refilled-idle sweeps, returning 429 + Retry-After + x-ratelimit-* headers and an S3-style XML body. - Key on trusted-proxy-validated ClientInfo.real_ip, else the socket peer address; never read spoofable X-Forwarded-For/X-Real-IP headers. Requests without a resolvable identity fail open. The echoed request id is charset-gated to prevent reflected XML injection. - Wire the layer once at startup via option_layer between CatchPanicLayer and ReadinessGateLayer (external stack only), gated by new RUSTFS_API_RATE_LIMIT_ENABLE/_RPM/_BURST constants; health and profiling probes, internode RPC/gRPC, and the console are exempt. - Make RUSTFS_CONSOLE_RATE_LIMIT_ENABLE/_RPM actually enforce by reusing the same limiter core through an axum middleware. - Delete the dead Swift ratelimit module, its isolated tests, and the stale logging-guardrail entry; keep the live SwiftError 429 mapping. - Add unit tests (exhaustion/recovery with injected time, concurrency, cap eviction, spoofed-header and fail-open behavior, env matrix) and e2e tests proving 429 + Retry-After on the real server and zero behavior change with default configuration.
This commit is contained in:
@@ -0,0 +1,52 @@
|
|||||||
|
// Copyright 2024 RustFS Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
/// Enable or disable per-client rate limiting for the S3 API.
|
||||||
|
///
|
||||||
|
/// When enabled (and `RUSTFS_API_RATE_LIMIT_RPM` > 0), requests are throttled
|
||||||
|
/// per client IP using a token bucket; over-limit requests receive
|
||||||
|
/// `429 Too Many Requests` with a `Retry-After` header. Internode RPC/gRPC,
|
||||||
|
/// health probes, and the console (which has its own limiter) are exempt.
|
||||||
|
/// Environment variable: RUSTFS_API_RATE_LIMIT_ENABLE
|
||||||
|
/// Example: RUSTFS_API_RATE_LIMIT_ENABLE=true
|
||||||
|
pub const ENV_API_RATE_LIMIT_ENABLE: &str = "RUSTFS_API_RATE_LIMIT_ENABLE";
|
||||||
|
|
||||||
|
/// Default for `RUSTFS_API_RATE_LIMIT_ENABLE`.
|
||||||
|
///
|
||||||
|
/// Disabled by default: RustFS ships permissive and operators opt in to
|
||||||
|
/// abuse-protection hardening. When disabled the request path is unchanged.
|
||||||
|
pub const DEFAULT_API_RATE_LIMIT_ENABLE: bool = false;
|
||||||
|
|
||||||
|
/// Sustained S3 API request budget per client IP, in requests per minute.
|
||||||
|
///
|
||||||
|
/// `0` means unlimited (rate limiting stays inert even when enabled).
|
||||||
|
/// Environment variable: RUSTFS_API_RATE_LIMIT_RPM
|
||||||
|
/// Example: RUSTFS_API_RATE_LIMIT_RPM=6000
|
||||||
|
pub const ENV_API_RATE_LIMIT_RPM: &str = "RUSTFS_API_RATE_LIMIT_RPM";
|
||||||
|
|
||||||
|
/// Default for `RUSTFS_API_RATE_LIMIT_RPM`.
|
||||||
|
///
|
||||||
|
/// `0` (unlimited) so that setting only the enable switch cannot throttle
|
||||||
|
/// traffic by surprise; operators must choose an explicit budget.
|
||||||
|
pub const DEFAULT_API_RATE_LIMIT_RPM: u32 = 0;
|
||||||
|
|
||||||
|
/// Burst capacity per client IP (maximum tokens in the bucket).
|
||||||
|
///
|
||||||
|
/// Allows short spikes above the sustained rate. `0` means "same as RPM".
|
||||||
|
/// Environment variable: RUSTFS_API_RATE_LIMIT_BURST
|
||||||
|
/// Example: RUSTFS_API_RATE_LIMIT_BURST=200
|
||||||
|
pub const ENV_API_RATE_LIMIT_BURST: &str = "RUSTFS_API_RATE_LIMIT_BURST";
|
||||||
|
|
||||||
|
/// Default for `RUSTFS_API_RATE_LIMIT_BURST` (`0` = same as RPM).
|
||||||
|
pub const DEFAULT_API_RATE_LIMIT_BURST: u32 = 0;
|
||||||
@@ -12,6 +12,7 @@
|
|||||||
// See the License for the specific language governing permissions and
|
// See the License for the specific language governing permissions and
|
||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
|
pub(crate) mod api;
|
||||||
pub(crate) mod app;
|
pub(crate) mod app;
|
||||||
pub(crate) mod body_limits;
|
pub(crate) mod body_limits;
|
||||||
pub(crate) mod capacity;
|
pub(crate) mod capacity;
|
||||||
|
|||||||
@@ -15,6 +15,8 @@
|
|||||||
#[cfg(feature = "constants")]
|
#[cfg(feature = "constants")]
|
||||||
pub mod constants;
|
pub mod constants;
|
||||||
#[cfg(feature = "constants")]
|
#[cfg(feature = "constants")]
|
||||||
|
pub use constants::api::*;
|
||||||
|
#[cfg(feature = "constants")]
|
||||||
pub use constants::app::*;
|
pub use constants::app::*;
|
||||||
#[cfg(feature = "constants")]
|
#[cfg(feature = "constants")]
|
||||||
pub use constants::body_limits::*;
|
pub use constants::body_limits::*;
|
||||||
|
|||||||
@@ -0,0 +1,109 @@
|
|||||||
|
// Copyright 2024 RustFS Team
|
||||||
|
//
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
|
||||||
|
//! E2E coverage for the opt-in per-client S3 API rate limit (backlog#1191):
|
||||||
|
//! the layer must be wired into the real server stack, reject over-limit
|
||||||
|
//! clients with `429` + `Retry-After`, keep health probes exempt, and stay
|
||||||
|
//! completely inert with default configuration.
|
||||||
|
|
||||||
|
use crate::common::{RustFSTestEnvironment, init_logging, local_http_client};
|
||||||
|
use serial_test::serial;
|
||||||
|
use tracing::info;
|
||||||
|
|
||||||
|
type TestResult = Result<(), Box<dyn std::error::Error + Send + Sync>>;
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn api_rate_limit_enforces_429_with_retry_after_when_enabled() -> TestResult {
|
||||||
|
init_logging();
|
||||||
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
|
// Burst 50 leaves headroom for the readiness-poll ListBuckets calls that
|
||||||
|
// share the loopback client IP; refill (60 rpm = 1/s) is slow enough that
|
||||||
|
// a rapid burst below reliably exhausts the bucket.
|
||||||
|
env.start_rustfs_server_with_env(
|
||||||
|
vec![],
|
||||||
|
&[
|
||||||
|
("RUSTFS_API_RATE_LIMIT_ENABLE", "true"),
|
||||||
|
("RUSTFS_API_RATE_LIMIT_RPM", "60"),
|
||||||
|
("RUSTFS_API_RATE_LIMIT_BURST", "50"),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
let client = local_http_client();
|
||||||
|
let list_buckets_url = format!("{}/", env.url);
|
||||||
|
|
||||||
|
let mut throttled = None;
|
||||||
|
let mut allowed = 0usize;
|
||||||
|
for _ in 0..80 {
|
||||||
|
let response = client.get(&list_buckets_url).send().await?;
|
||||||
|
if response.status() == reqwest::StatusCode::TOO_MANY_REQUESTS {
|
||||||
|
throttled = Some(response);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
allowed += 1;
|
||||||
|
}
|
||||||
|
|
||||||
|
let throttled = throttled.unwrap_or_else(|| panic!("no 429 within 80 rapid requests ({allowed} allowed) at burst 50"));
|
||||||
|
assert!(allowed > 0, "healthy traffic below the burst must not be throttled");
|
||||||
|
info!("rate limit engaged after {allowed} allowed requests");
|
||||||
|
|
||||||
|
let retry_after = throttled
|
||||||
|
.headers()
|
||||||
|
.get(reqwest::header::RETRY_AFTER)
|
||||||
|
.and_then(|v| v.to_str().ok())
|
||||||
|
.and_then(|v| v.parse::<u64>().ok())
|
||||||
|
.expect("429 must carry a numeric Retry-After header");
|
||||||
|
assert!(retry_after >= 1, "Retry-After must be at least one second, got {retry_after}");
|
||||||
|
assert_eq!(
|
||||||
|
throttled.headers().get("x-ratelimit-limit").and_then(|v| v.to_str().ok()),
|
||||||
|
Some("60"),
|
||||||
|
"429 must expose the configured limit"
|
||||||
|
);
|
||||||
|
|
||||||
|
let body = throttled.text().await?;
|
||||||
|
assert!(body.contains("<Code>TooManyRequests</Code>"), "S3-style error body expected: {body}");
|
||||||
|
|
||||||
|
// Health probes stay exempt even while the client budget is exhausted.
|
||||||
|
let health = client.get(format!("{}/health", env.url)).send().await?;
|
||||||
|
assert_ne!(
|
||||||
|
health.status(),
|
||||||
|
reqwest::StatusCode::TOO_MANY_REQUESTS,
|
||||||
|
"health probes must never be rate limited"
|
||||||
|
);
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn api_rate_limit_stays_inert_by_default() -> TestResult {
|
||||||
|
init_logging();
|
||||||
|
let mut env = RustFSTestEnvironment::new().await?;
|
||||||
|
env.start_rustfs_server(vec![]).await?;
|
||||||
|
|
||||||
|
let client = local_http_client();
|
||||||
|
let list_buckets_url = format!("{}/", env.url);
|
||||||
|
|
||||||
|
for i in 0..80 {
|
||||||
|
let response = client.get(&list_buckets_url).send().await?;
|
||||||
|
assert_ne!(
|
||||||
|
response.status(),
|
||||||
|
reqwest::StatusCode::TOO_MANY_REQUESTS,
|
||||||
|
"request {i} was throttled although rate limiting is disabled by default"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
@@ -67,6 +67,10 @@ mod bucket_policy_check_test;
|
|||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod security_boundary_test;
|
mod security_boundary_test;
|
||||||
|
|
||||||
|
// Opt-in per-client S3 API rate limiting (backlog#1191)
|
||||||
|
#[cfg(test)]
|
||||||
|
mod api_rate_limit_test;
|
||||||
|
|
||||||
// Admin authorization gate: non-admin denial + root-credential lifecycle (backlog#1151 sec-4)
|
// Admin authorization gate: non-admin denial + root-credential lifecycle (backlog#1151 sec-4)
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod admin_auth_test;
|
mod admin_auth_test;
|
||||||
|
|||||||
@@ -46,7 +46,6 @@ pub mod formpost;
|
|||||||
pub mod handler;
|
pub mod handler;
|
||||||
pub mod object;
|
pub mod object;
|
||||||
pub mod quota;
|
pub mod quota;
|
||||||
pub mod ratelimit;
|
|
||||||
pub mod router;
|
pub mod router;
|
||||||
pub mod slo;
|
pub mod slo;
|
||||||
pub mod staticweb;
|
pub mod staticweb;
|
||||||
|
|||||||
@@ -1,434 +0,0 @@
|
|||||||
// Copyright 2024 RustFS Team
|
|
||||||
//
|
|
||||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
|
||||||
// you may not use this file except in compliance with the License.
|
|
||||||
// You may obtain a copy of the License at
|
|
||||||
//
|
|
||||||
// http://www.apache.org/licenses/LICENSE-2.0
|
|
||||||
//
|
|
||||||
// Unless required by applicable law or agreed to in writing, software
|
|
||||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
|
||||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
||||||
// See the License for the specific language governing permissions and
|
|
||||||
// limitations under the License.
|
|
||||||
|
|
||||||
//! Rate Limiting Support for Swift API
|
|
||||||
//!
|
|
||||||
//! This module implements rate limiting to prevent abuse and ensure fair resource
|
|
||||||
//! allocation across tenants. Rate limits can be applied per-account, per-container,
|
|
||||||
//! or per-IP address.
|
|
||||||
//!
|
|
||||||
//! # Configuration
|
|
||||||
//!
|
|
||||||
//! Rate limits are configured via container metadata:
|
|
||||||
//!
|
|
||||||
//! ```bash
|
|
||||||
//! # Set account-level rate limit: 1000 requests per minute
|
|
||||||
//! swift post -m "X-Account-Meta-Rate-Limit:1000/60"
|
|
||||||
//!
|
|
||||||
//! # Set container-level rate limit: 100 requests per minute
|
|
||||||
//! swift post container -m "X-Container-Meta-Rate-Limit:100/60"
|
|
||||||
//! ```
|
|
||||||
//!
|
|
||||||
//! # Response Headers
|
|
||||||
//!
|
|
||||||
//! Rate limit information is included in all responses:
|
|
||||||
//!
|
|
||||||
//! ```http
|
|
||||||
//! HTTP/1.1 200 OK
|
|
||||||
//! X-RateLimit-Limit: 1000
|
|
||||||
//! X-RateLimit-Remaining: 950
|
|
||||||
//! X-RateLimit-Reset: 1740003600
|
|
||||||
//! ```
|
|
||||||
//!
|
|
||||||
//! When rate limit is exceeded:
|
|
||||||
//!
|
|
||||||
//! ```http
|
|
||||||
//! HTTP/1.1 429 Too Many Requests
|
|
||||||
//! X-RateLimit-Limit: 1000
|
|
||||||
//! X-RateLimit-Remaining: 0
|
|
||||||
//! X-RateLimit-Reset: 1740003600
|
|
||||||
//! Retry-After: 30
|
|
||||||
//! ```
|
|
||||||
//!
|
|
||||||
//! # Algorithm
|
|
||||||
//!
|
|
||||||
//! Uses token bucket algorithm with per-second refill rate:
|
|
||||||
//! - Each request consumes 1 token
|
|
||||||
//! - Tokens refill at configured rate
|
|
||||||
//! - Burst capacity allows temporary spikes
|
|
||||||
|
|
||||||
use super::{SwiftError, SwiftResult};
|
|
||||||
use std::collections::HashMap;
|
|
||||||
use std::sync::{Arc, Mutex};
|
|
||||||
use std::time::{SystemTime, UNIX_EPOCH};
|
|
||||||
use tracing::debug;
|
|
||||||
|
|
||||||
const LOG_COMPONENT_PROTOCOLS: &str = "protocols";
|
|
||||||
const LOG_SUBSYSTEM_SWIFT_RATELIMIT: &str = "swift_ratelimit";
|
|
||||||
const EVENT_SWIFT_RATELIMIT_STATE: &str = "swift_ratelimit_state";
|
|
||||||
|
|
||||||
/// Rate limit configuration
|
|
||||||
#[derive(Debug, Clone, PartialEq)]
|
|
||||||
pub struct RateLimit {
|
|
||||||
/// Maximum requests allowed in time window
|
|
||||||
pub limit: u32,
|
|
||||||
|
|
||||||
/// Time window in seconds
|
|
||||||
pub window_seconds: u32,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl RateLimit {
|
|
||||||
/// Parse rate limit from metadata value
|
|
||||||
///
|
|
||||||
/// Format: "limit/window_seconds" (e.g., "1000/60" = 1000 requests per 60 seconds)
|
|
||||||
pub fn parse(value: &str) -> SwiftResult<Self> {
|
|
||||||
let parts: Vec<&str> = value.split('/').collect();
|
|
||||||
if parts.len() != 2 {
|
|
||||||
return Err(SwiftError::BadRequest(format!(
|
|
||||||
"Invalid rate limit format: {}. Expected format: limit/window_seconds",
|
|
||||||
value
|
|
||||||
)));
|
|
||||||
}
|
|
||||||
|
|
||||||
let limit = parts[0]
|
|
||||||
.parse::<u32>()
|
|
||||||
.map_err(|_| SwiftError::BadRequest(format!("Invalid rate limit value: {}", parts[0])))?;
|
|
||||||
|
|
||||||
let window_seconds = parts[1]
|
|
||||||
.parse::<u32>()
|
|
||||||
.map_err(|_| SwiftError::BadRequest(format!("Invalid window value: {}", parts[1])))?;
|
|
||||||
|
|
||||||
if window_seconds == 0 {
|
|
||||||
return Err(SwiftError::BadRequest("Rate limit window cannot be zero".to_string()));
|
|
||||||
}
|
|
||||||
|
|
||||||
Ok(RateLimit { limit, window_seconds })
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Calculate refill rate (tokens per second)
|
|
||||||
pub fn refill_rate(&self) -> f64 {
|
|
||||||
self.limit as f64 / self.window_seconds as f64
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Token bucket for rate limiting
|
|
||||||
#[derive(Debug, Clone)]
|
|
||||||
struct TokenBucket {
|
|
||||||
/// Maximum tokens (burst capacity)
|
|
||||||
capacity: u32,
|
|
||||||
|
|
||||||
/// Current available tokens
|
|
||||||
tokens: f64,
|
|
||||||
|
|
||||||
/// Refill rate (tokens per second)
|
|
||||||
refill_rate: f64,
|
|
||||||
|
|
||||||
/// Last refill timestamp (Unix seconds)
|
|
||||||
last_refill: u64,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl TokenBucket {
|
|
||||||
fn new(rate_limit: &RateLimit) -> Self {
|
|
||||||
let capacity = rate_limit.limit;
|
|
||||||
let refill_rate = rate_limit.refill_rate();
|
|
||||||
|
|
||||||
TokenBucket {
|
|
||||||
capacity,
|
|
||||||
tokens: capacity as f64, // Start full
|
|
||||||
refill_rate,
|
|
||||||
last_refill: current_timestamp(),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Try to consume a token
|
|
||||||
///
|
|
||||||
/// Returns Ok(remaining_tokens) if successful, Err(retry_after_seconds) if rate limited
|
|
||||||
fn try_consume(&mut self) -> Result<u32, u64> {
|
|
||||||
// Refill tokens based on time elapsed
|
|
||||||
let now = current_timestamp();
|
|
||||||
let elapsed = now.saturating_sub(self.last_refill);
|
|
||||||
|
|
||||||
if elapsed > 0 {
|
|
||||||
let refill_amount = self.refill_rate * elapsed as f64;
|
|
||||||
self.tokens = (self.tokens + refill_amount).min(self.capacity as f64);
|
|
||||||
self.last_refill = now;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Try to consume 1 token
|
|
||||||
if self.tokens >= 1.0 {
|
|
||||||
self.tokens -= 1.0;
|
|
||||||
Ok(self.tokens.floor() as u32)
|
|
||||||
} else {
|
|
||||||
// Calculate retry-after: time until 1 token is available
|
|
||||||
let tokens_needed = 1.0 - self.tokens;
|
|
||||||
let retry_after = (tokens_needed / self.refill_rate).ceil() as u64;
|
|
||||||
Err(retry_after)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get current token count
|
|
||||||
fn remaining(&mut self) -> u32 {
|
|
||||||
// Refill tokens based on time elapsed
|
|
||||||
let now = current_timestamp();
|
|
||||||
let elapsed = now.saturating_sub(self.last_refill);
|
|
||||||
|
|
||||||
if elapsed > 0 {
|
|
||||||
let refill_amount = self.refill_rate * elapsed as f64;
|
|
||||||
self.tokens = (self.tokens + refill_amount).min(self.capacity as f64);
|
|
||||||
self.last_refill = now;
|
|
||||||
}
|
|
||||||
|
|
||||||
self.tokens.floor() as u32
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get reset timestamp (when bucket will be full)
|
|
||||||
fn reset_timestamp(&self, now: u64) -> u64 {
|
|
||||||
if self.tokens >= self.capacity as f64 {
|
|
||||||
now
|
|
||||||
} else {
|
|
||||||
let tokens_to_refill = self.capacity as f64 - self.tokens;
|
|
||||||
let seconds_to_full = (tokens_to_refill / self.refill_rate).ceil() as u64;
|
|
||||||
now + seconds_to_full
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Global rate limiter state (in-memory)
|
|
||||||
///
|
|
||||||
/// In production, this should be backed by Redis or similar distributed store
|
|
||||||
#[derive(Clone)]
|
|
||||||
pub struct RateLimiter {
|
|
||||||
buckets: Arc<Mutex<HashMap<String, TokenBucket>>>,
|
|
||||||
}
|
|
||||||
|
|
||||||
impl RateLimiter {
|
|
||||||
/// Create new rate limiter
|
|
||||||
pub fn new() -> Self {
|
|
||||||
RateLimiter {
|
|
||||||
buckets: Arc::new(Mutex::new(HashMap::new())),
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Check and consume rate limit quota
|
|
||||||
///
|
|
||||||
/// Returns (remaining, reset_timestamp) if successful,
|
|
||||||
/// or SwiftError::TooManyRequests if rate limited
|
|
||||||
pub fn check_rate_limit(&self, key: &str, rate_limit: &RateLimit) -> SwiftResult<(u32, u64)> {
|
|
||||||
let mut buckets = self.buckets.lock().expect("operation should succeed");
|
|
||||||
|
|
||||||
// Get or create bucket for this key
|
|
||||||
let bucket = buckets.entry(key.to_string()).or_insert_with(|| TokenBucket::new(rate_limit));
|
|
||||||
|
|
||||||
let now = current_timestamp();
|
|
||||||
let reset = bucket.reset_timestamp(now);
|
|
||||||
|
|
||||||
match bucket.try_consume() {
|
|
||||||
Ok(remaining) => {
|
|
||||||
debug!(event = EVENT_SWIFT_RATELIMIT_STATE, component = LOG_COMPONENT_PROTOCOLS, subsystem = LOG_SUBSYSTEM_SWIFT_RATELIMIT, key = %key, remaining, result = "allowed", "swift ratelimit state changed");
|
|
||||||
Ok((remaining, reset))
|
|
||||||
}
|
|
||||||
Err(retry_after) => {
|
|
||||||
debug!(event = EVENT_SWIFT_RATELIMIT_STATE, component = LOG_COMPONENT_PROTOCOLS, subsystem = LOG_SUBSYSTEM_SWIFT_RATELIMIT, key = %key, retry_after, result = "limited", "swift ratelimit state changed");
|
|
||||||
Err(SwiftError::TooManyRequests {
|
|
||||||
retry_after,
|
|
||||||
limit: rate_limit.limit,
|
|
||||||
reset,
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get current rate limit status without consuming quota
|
|
||||||
pub fn get_status(&self, key: &str, rate_limit: &RateLimit) -> (u32, u64) {
|
|
||||||
let mut buckets = self.buckets.lock().expect("operation should succeed");
|
|
||||||
|
|
||||||
let bucket = buckets.entry(key.to_string()).or_insert_with(|| TokenBucket::new(rate_limit));
|
|
||||||
|
|
||||||
let now = current_timestamp();
|
|
||||||
let remaining = bucket.remaining();
|
|
||||||
let reset = bucket.reset_timestamp(now);
|
|
||||||
|
|
||||||
(remaining, reset)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
impl Default for RateLimiter {
|
|
||||||
fn default() -> Self {
|
|
||||||
Self::new()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Get current Unix timestamp in seconds
|
|
||||||
fn current_timestamp() -> u64 {
|
|
||||||
SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Extract rate limit from account or container metadata
|
|
||||||
pub fn extract_rate_limit(metadata: &HashMap<String, String>) -> Option<RateLimit> {
|
|
||||||
// Check for rate limit in metadata
|
|
||||||
if let Some(rate_limit_str) = metadata
|
|
||||||
.get("x-account-meta-rate-limit")
|
|
||||||
.or_else(|| metadata.get("x-container-meta-rate-limit"))
|
|
||||||
{
|
|
||||||
RateLimit::parse(rate_limit_str).ok()
|
|
||||||
} else {
|
|
||||||
None
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Build rate limit key for tracking
|
|
||||||
pub fn build_rate_limit_key(account: &str, container: Option<&str>) -> String {
|
|
||||||
if let Some(cont) = container {
|
|
||||||
format!("account:{}:container:{}", account, cont)
|
|
||||||
} else {
|
|
||||||
format!("account:{}", account)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
#[cfg(test)]
|
|
||||||
mod tests {
|
|
||||||
use super::*;
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_parse_rate_limit_valid() {
|
|
||||||
let rate_limit = RateLimit::parse("1000/60").expect("operation should succeed");
|
|
||||||
assert_eq!(rate_limit.limit, 1000);
|
|
||||||
assert_eq!(rate_limit.window_seconds, 60);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_parse_rate_limit_invalid_format() {
|
|
||||||
let result = RateLimit::parse("1000");
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_parse_rate_limit_invalid_limit() {
|
|
||||||
let result = RateLimit::parse("not_a_number/60");
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_parse_rate_limit_invalid_window() {
|
|
||||||
let result = RateLimit::parse("1000/not_a_number");
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_parse_rate_limit_zero_window() {
|
|
||||||
let result = RateLimit::parse("1000/0");
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_rate_limit_refill_rate() {
|
|
||||||
let rate_limit = RateLimit {
|
|
||||||
limit: 1000,
|
|
||||||
window_seconds: 60,
|
|
||||||
};
|
|
||||||
assert!((rate_limit.refill_rate() - 16.666666).abs() < 0.001);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_token_bucket_consume() {
|
|
||||||
let rate_limit = RateLimit {
|
|
||||||
limit: 10,
|
|
||||||
window_seconds: 60,
|
|
||||||
};
|
|
||||||
let mut bucket = TokenBucket::new(&rate_limit);
|
|
||||||
|
|
||||||
// Should be able to consume up to limit
|
|
||||||
for i in 0..10 {
|
|
||||||
let result = bucket.try_consume();
|
|
||||||
assert!(result.is_ok(), "Token {} should succeed", i);
|
|
||||||
}
|
|
||||||
|
|
||||||
// 11th request should fail
|
|
||||||
let result = bucket.try_consume();
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_token_bucket_remaining() {
|
|
||||||
let rate_limit = RateLimit {
|
|
||||||
limit: 100,
|
|
||||||
window_seconds: 60,
|
|
||||||
};
|
|
||||||
let mut bucket = TokenBucket::new(&rate_limit);
|
|
||||||
|
|
||||||
// Initial: 100 tokens
|
|
||||||
assert_eq!(bucket.remaining(), 100);
|
|
||||||
|
|
||||||
// Consume 10
|
|
||||||
for _ in 0..10 {
|
|
||||||
bucket.try_consume().expect("operation should succeed");
|
|
||||||
}
|
|
||||||
|
|
||||||
assert_eq!(bucket.remaining(), 90);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_rate_limiter() {
|
|
||||||
let limiter = RateLimiter::new();
|
|
||||||
let rate_limit = RateLimit {
|
|
||||||
limit: 5,
|
|
||||||
window_seconds: 60,
|
|
||||||
};
|
|
||||||
|
|
||||||
// Should allow 5 requests
|
|
||||||
for _ in 0..5 {
|
|
||||||
let result = limiter.check_rate_limit("test_key", &rate_limit);
|
|
||||||
assert!(result.is_ok());
|
|
||||||
}
|
|
||||||
|
|
||||||
// 6th request should fail
|
|
||||||
let result = limiter.check_rate_limit("test_key", &rate_limit);
|
|
||||||
assert!(result.is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_extract_rate_limit_account() {
|
|
||||||
let mut metadata = HashMap::new();
|
|
||||||
metadata.insert("x-account-meta-rate-limit".to_string(), "1000/60".to_string());
|
|
||||||
|
|
||||||
let rate_limit = extract_rate_limit(&metadata);
|
|
||||||
assert!(rate_limit.is_some());
|
|
||||||
|
|
||||||
let rate_limit = rate_limit.expect("operation should succeed");
|
|
||||||
assert_eq!(rate_limit.limit, 1000);
|
|
||||||
assert_eq!(rate_limit.window_seconds, 60);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_extract_rate_limit_container() {
|
|
||||||
let mut metadata = HashMap::new();
|
|
||||||
metadata.insert("x-container-meta-rate-limit".to_string(), "100/60".to_string());
|
|
||||||
|
|
||||||
let rate_limit = extract_rate_limit(&metadata);
|
|
||||||
assert!(rate_limit.is_some());
|
|
||||||
|
|
||||||
let rate_limit = rate_limit.expect("operation should succeed");
|
|
||||||
assert_eq!(rate_limit.limit, 100);
|
|
||||||
assert_eq!(rate_limit.window_seconds, 60);
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_extract_rate_limit_none() {
|
|
||||||
let metadata = HashMap::new();
|
|
||||||
let rate_limit = extract_rate_limit(&metadata);
|
|
||||||
assert!(rate_limit.is_none());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_build_rate_limit_key_account() {
|
|
||||||
let key = build_rate_limit_key("AUTH_test", None);
|
|
||||||
assert_eq!(key, "account:AUTH_test");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_build_rate_limit_key_container() {
|
|
||||||
let key = build_rate_limit_key("AUTH_test", Some("my-container"));
|
|
||||||
assert_eq!(key, "account:AUTH_test:container:my-container");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -43,36 +43,4 @@ mod swift_integration {
|
|||||||
let parsed = expiration::parse_delete_at(delete_at).unwrap();
|
let parsed = expiration::parse_delete_at(delete_at).unwrap();
|
||||||
assert_eq!(parsed, 1740000000);
|
assert_eq!(parsed, 1740000000);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_multiple_rate_limit_keys() {
|
|
||||||
let limiter = ratelimit::RateLimiter::new();
|
|
||||||
let rate = ratelimit::RateLimit {
|
|
||||||
limit: 3,
|
|
||||||
window_seconds: 60,
|
|
||||||
};
|
|
||||||
|
|
||||||
// Different keys should have separate limits
|
|
||||||
for _ in 0..3 {
|
|
||||||
assert!(limiter.check_rate_limit("key1", &rate).is_ok());
|
|
||||||
assert!(limiter.check_rate_limit("key2", &rate).is_ok());
|
|
||||||
}
|
|
||||||
|
|
||||||
// Both keys should now be exhausted
|
|
||||||
assert!(limiter.check_rate_limit("key1", &rate).is_err());
|
|
||||||
assert!(limiter.check_rate_limit("key2", &rate).is_err());
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
|
||||||
fn test_rate_limit_metadata_extraction() {
|
|
||||||
let mut metadata = HashMap::new();
|
|
||||||
metadata.insert("x-account-meta-rate-limit".to_string(), "1000/60".to_string());
|
|
||||||
|
|
||||||
let rate_limit = ratelimit::extract_rate_limit(&metadata);
|
|
||||||
assert!(rate_limit.is_some());
|
|
||||||
|
|
||||||
let rate_limit = rate_limit.unwrap();
|
|
||||||
assert_eq!(rate_limit.limit, 1000);
|
|
||||||
assert_eq!(rate_limit.window_seconds, 60);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -16,7 +16,7 @@
|
|||||||
|
|
||||||
#![cfg(feature = "swift")]
|
#![cfg(feature = "swift")]
|
||||||
|
|
||||||
use rustfs_protocols::swift::{quota, ratelimit, slo, symlink, sync, tempurl, versioning};
|
use rustfs_protocols::swift::{quota, slo, symlink, sync, tempurl, versioning};
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
|
|
||||||
/// Test sync configuration parsing
|
/// Test sync configuration parsing
|
||||||
@@ -92,14 +92,6 @@ fn test_symlink_detection() {
|
|||||||
let _is_symlink = symlink::is_symlink(&metadata);
|
let _is_symlink = symlink::is_symlink(&metadata);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Test rate limit parsing
|
|
||||||
#[test]
|
|
||||||
fn test_rate_limit_parsing() {
|
|
||||||
let rl = ratelimit::RateLimit::parse("100/60").unwrap();
|
|
||||||
assert_eq!(rl.limit, 100);
|
|
||||||
assert_eq!(rl.window_seconds, 60);
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Test quota structure
|
/// Test quota structure
|
||||||
#[test]
|
#[test]
|
||||||
fn test_quota_structure() {
|
fn test_quota_structure() {
|
||||||
|
|||||||
@@ -16,6 +16,10 @@ use crate::admin::runtime_sources::{current_oidc_handle, default_admin_usecase};
|
|||||||
use crate::admin::storage_api::access::RequestContext;
|
use crate::admin::storage_api::access::RequestContext;
|
||||||
use crate::license::has_valid_license;
|
use crate::license::has_valid_license;
|
||||||
use crate::server::has_path_prefix;
|
use crate::server::has_path_prefix;
|
||||||
|
use crate::server::rate_limit::{
|
||||||
|
LABEL_RATE_LIMIT_SCOPE, METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL, RATE_LIMIT_SCOPE_CONSOLE, RateLimitDecision,
|
||||||
|
RateLimitQuota, RateLimiter, apply_throttle_headers, client_ip,
|
||||||
|
};
|
||||||
use crate::server::{
|
use crate::server::{
|
||||||
CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HeaderMapCarrier, HealthProbe, LICENSE, RUSTFS_ADMIN_PREFIX,
|
CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HeaderMapCarrier, HealthProbe, LICENSE, RUSTFS_ADMIN_PREFIX,
|
||||||
RequestContextLayer, VERSION, build_health_response_parts, collect_probe_readiness,
|
RequestContextLayer, VERSION, build_health_response_parts, collect_probe_readiness,
|
||||||
@@ -36,7 +40,7 @@ use rust_embed::RustEmbed;
|
|||||||
use serde::Serialize;
|
use serde::Serialize;
|
||||||
use std::{
|
use std::{
|
||||||
net::{IpAddr, SocketAddr},
|
net::{IpAddr, SocketAddr},
|
||||||
sync::OnceLock,
|
sync::{Arc, OnceLock},
|
||||||
time::Duration,
|
time::Duration,
|
||||||
};
|
};
|
||||||
use tower_http::catch_panic::CatchPanicLayer;
|
use tower_http::catch_panic::CatchPanicLayer;
|
||||||
@@ -569,12 +573,38 @@ fn setup_console_middleware_stack(
|
|||||||
// Add request body limit (10MB for console uploads)
|
// Add request body limit (10MB for console uploads)
|
||||||
.layer(RequestBodyLimitLayer::new(5 * 1024 * 1024 * 1024));
|
.layer(RequestBodyLimitLayer::new(5 * 1024 * 1024 * 1024));
|
||||||
|
|
||||||
// Add rate limiting if enabled
|
// Per-client console rate limiting, sharing the S3 API limiter core
|
||||||
|
// (`server::rate_limit`). Added last so it is the outermost layer and
|
||||||
|
// over-limit requests are rejected before any other middleware runs.
|
||||||
|
// The client identity comes from validated request extensions only
|
||||||
|
// (trusted-proxy `ClientInfo`, then socket peer address); requests
|
||||||
|
// without one fail open rather than trusting spoofable headers.
|
||||||
if rate_limit_enable {
|
if rate_limit_enable {
|
||||||
info!("Console rate limiting enabled: {} requests per minute", rate_limit_rpm);
|
if let Some(quota) = RateLimitQuota::per_minute(rate_limit_rpm, 0) {
|
||||||
// Note: tower-http doesn't provide a built-in rate limiter, but we have the foundation
|
let limiter = Arc::new(RateLimiter::new(quota));
|
||||||
// For production, you would integrate with a rate limiting service like Redis
|
app = app.layer(middleware::from_fn(move |req: Request, next: middleware::Next| {
|
||||||
// For now, we log that it's configured and ready for integration
|
let limiter = limiter.clone();
|
||||||
|
async move {
|
||||||
|
match client_ip(&req).map(|ip| limiter.check(ip)) {
|
||||||
|
Some(RateLimitDecision::Limited(throttle)) => {
|
||||||
|
metrics::counter!(
|
||||||
|
METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL,
|
||||||
|
LABEL_RATE_LIMIT_SCOPE => RATE_LIMIT_SCOPE_CONSOLE
|
||||||
|
)
|
||||||
|
.increment(1);
|
||||||
|
debug!(retry_after_secs = throttle.retry_after_secs, "Console request rejected by rate limit");
|
||||||
|
let mut response = StatusCode::TOO_MANY_REQUESTS.into_response();
|
||||||
|
apply_throttle_headers(response.headers_mut(), limiter.quota().requests_per_minute, &throttle);
|
||||||
|
response
|
||||||
|
}
|
||||||
|
_ => next.run(req).await,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}));
|
||||||
|
info!("Console rate limiting enabled: {} requests per minute", rate_limit_rpm);
|
||||||
|
} else {
|
||||||
|
warn!("Console rate limiting enabled but RPM is 0; it stays inactive");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
app
|
app
|
||||||
@@ -864,6 +894,39 @@ mod tests {
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// setup_console_middleware_stack reads ENV_HEALTH_ENDPOINT_ENABLE (see above).
|
||||||
|
#[tokio::test]
|
||||||
|
#[serial]
|
||||||
|
async fn console_rate_limit_rejects_over_limit_clients_with_429() {
|
||||||
|
let app = setup_console_middleware_stack(parse_cors_origins(None), true, 1, 30);
|
||||||
|
|
||||||
|
let request_for = |last_octet: u8| {
|
||||||
|
let mut request = Request::builder()
|
||||||
|
.uri(format!("{CONSOLE_PREFIX}/index.html"))
|
||||||
|
.body(Body::empty())
|
||||||
|
.expect("failed to build request");
|
||||||
|
request.extensions_mut().insert(crate::server::RemoteAddr(SocketAddr::new(
|
||||||
|
IpAddr::from([198, 51, 100, last_octet]),
|
||||||
|
50000,
|
||||||
|
)));
|
||||||
|
request
|
||||||
|
};
|
||||||
|
|
||||||
|
let first = app.clone().oneshot(request_for(7)).await.expect("request should succeed");
|
||||||
|
assert_ne!(first.status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
|
||||||
|
let second = app.clone().oneshot(request_for(7)).await.expect("request should succeed");
|
||||||
|
assert_eq!(second.status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
assert!(
|
||||||
|
second.headers().contains_key(http::header::RETRY_AFTER),
|
||||||
|
"console 429 must carry Retry-After"
|
||||||
|
);
|
||||||
|
|
||||||
|
// A different client IP has its own budget.
|
||||||
|
let other = app.oneshot(request_for(8)).await.expect("request should succeed");
|
||||||
|
assert_ne!(other.status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test(flavor = "current_thread")]
|
#[tokio::test(flavor = "current_thread")]
|
||||||
async fn console_trace_layer_records_request_id_on_current_span() {
|
async fn console_trace_layer_records_request_id_on_current_span() {
|
||||||
let writer = SharedWriter::default();
|
let writer = SharedWriter::default();
|
||||||
|
|||||||
+59
-17
@@ -27,6 +27,7 @@ use crate::server::{
|
|||||||
RedirectLayer, RequestContextLayer, RequestLoggingLayer, S3ErrorMessageCompatLayer, VirtualHostStyleHintLayer,
|
RedirectLayer, RequestContextLayer, RequestLoggingLayer, S3ErrorMessageCompatLayer, VirtualHostStyleHintLayer,
|
||||||
redact_sensitive_uri_query,
|
redact_sensitive_uri_query,
|
||||||
},
|
},
|
||||||
|
rate_limit::{RateLimitLayer, api_rate_limit_layer_from_env},
|
||||||
tls_material::{
|
tls_material::{
|
||||||
TlsAcceptorHolder, TlsHandshakeFailureKind, build_acceptor_from_loaded, load_tls_material, spawn_reload_loop,
|
TlsAcceptorHolder, TlsHandshakeFailureKind, build_acceptor_from_loaded, load_tls_material, spawn_reload_loop,
|
||||||
},
|
},
|
||||||
@@ -111,6 +112,7 @@ const EVENT_HTTP_BIND_FAILED: &str = "http_bind_failed";
|
|||||||
const EVENT_HTTP_STARTUP_ENDPOINTS: &str = "http_startup_endpoints";
|
const EVENT_HTTP_STARTUP_ENDPOINTS: &str = "http_startup_endpoints";
|
||||||
const EVENT_HTTP_HOST_ROUTING: &str = "http_host_routing";
|
const EVENT_HTTP_HOST_ROUTING: &str = "http_host_routing";
|
||||||
const EVENT_HTTP_COMPRESSION_STATE: &str = "http_compression_state";
|
const EVENT_HTTP_COMPRESSION_STATE: &str = "http_compression_state";
|
||||||
|
const EVENT_API_RATE_LIMIT_STATE: &str = "api_rate_limit_state";
|
||||||
const EVENT_HTTP_TRANSPORT_PARAMETERS: &str = "http_transport_parameters";
|
const EVENT_HTTP_TRANSPORT_PARAMETERS: &str = "http_transport_parameters";
|
||||||
const EVENT_HTTP_ACCEPT_LOOP_STATE: &str = "http_accept_loop_state";
|
const EVENT_HTTP_ACCEPT_LOOP_STATE: &str = "http_accept_loop_state";
|
||||||
const EVENT_HTTP_CONNECTION_DRAIN: &str = "http_connection_drain";
|
const EVENT_HTTP_CONNECTION_DRAIN: &str = "http_connection_drain";
|
||||||
@@ -700,6 +702,34 @@ pub async fn start_http_server(
|
|||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Per-client S3 API rate limiting (backlog#1191). Built once so every
|
||||||
|
// connection's stack shares the same limiter state; `None` (the default)
|
||||||
|
// leaves the request path unchanged.
|
||||||
|
let api_rate_limit_layer = api_rate_limit_layer_from_env();
|
||||||
|
match &api_rate_limit_layer {
|
||||||
|
Some(layer) => {
|
||||||
|
let quota = layer.quota();
|
||||||
|
info!(
|
||||||
|
event = EVENT_API_RATE_LIMIT_STATE,
|
||||||
|
component = LOG_COMPONENT_SERVER,
|
||||||
|
subsystem = LOG_SUBSYSTEM_HTTP,
|
||||||
|
state = "enabled",
|
||||||
|
requests_per_minute = quota.requests_per_minute,
|
||||||
|
burst = quota.burst,
|
||||||
|
"API rate limit state changed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
None => {
|
||||||
|
debug!(
|
||||||
|
event = EVENT_API_RATE_LIMIT_STATE,
|
||||||
|
component = LOG_COMPONENT_SERVER,
|
||||||
|
subsystem = LOG_SUBSYSTEM_HTTP,
|
||||||
|
state = "disabled",
|
||||||
|
"API rate limit state changed"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
let is_console = config.console_enable;
|
let is_console = config.console_enable;
|
||||||
let server_domains_configured = !config.server_domains.is_empty();
|
let server_domains_configured = !config.server_domains.is_empty();
|
||||||
let task_handle = tokio::spawn(async move {
|
let task_handle = tokio::spawn(async move {
|
||||||
@@ -890,6 +920,7 @@ pub async fn start_http_server(
|
|||||||
readiness: readiness.clone(),
|
readiness: readiness.clone(),
|
||||||
keystone_auth: auth_keystone::get_keystone_auth(),
|
keystone_auth: auth_keystone::get_keystone_auth(),
|
||||||
trusted_proxy_layer: rustfs_trusted_proxies::is_enabled().then(|| rustfs_trusted_proxies::layer().clone()),
|
trusted_proxy_layer: rustfs_trusted_proxies::is_enabled().then(|| rustfs_trusted_proxies::layer().clone()),
|
||||||
|
rate_limit_layer: api_rate_limit_layer.clone(),
|
||||||
};
|
};
|
||||||
|
|
||||||
process_connection(socket, tls_acceptor.clone(), connection_ctx, graceful.watcher());
|
process_connection(socket, tls_acceptor.clone(), connection_ctx, graceful.watcher());
|
||||||
@@ -946,6 +977,9 @@ struct ConnectionContext {
|
|||||||
keystone_auth: Option<Arc<rustfs_keystone::KeystoneAuthProvider>>,
|
keystone_auth: Option<Arc<rustfs_keystone::KeystoneAuthProvider>>,
|
||||||
/// Pre-computed trusted proxy layer (avoids per-connection is_enabled() check).
|
/// Pre-computed trusted proxy layer (avoids per-connection is_enabled() check).
|
||||||
trusted_proxy_layer: Option<rustfs_trusted_proxies::TrustedProxyLayer>,
|
trusted_proxy_layer: Option<rustfs_trusted_proxies::TrustedProxyLayer>,
|
||||||
|
/// Per-client API rate limit layer; `None` when disabled (the default).
|
||||||
|
/// All clones share one limiter, keeping budgets global across connections.
|
||||||
|
rate_limit_layer: Option<RateLimitLayer>,
|
||||||
}
|
}
|
||||||
|
|
||||||
#[derive(Clone)]
|
#[derive(Clone)]
|
||||||
@@ -1130,6 +1164,7 @@ fn process_connection(
|
|||||||
readiness,
|
readiness,
|
||||||
keystone_auth,
|
keystone_auth,
|
||||||
trusted_proxy_layer,
|
trusted_proxy_layer,
|
||||||
|
rate_limit_layer,
|
||||||
} = context;
|
} = context;
|
||||||
|
|
||||||
// Build the hybrid service per-connection.
|
// Build the hybrid service per-connection.
|
||||||
@@ -1183,23 +1218,24 @@ fn process_connection(
|
|||||||
// 5. RequestContextLayer — creates RequestContext in extensions
|
// 5. RequestContextLayer — creates RequestContext in extensions
|
||||||
// 6. EmptyBodyContentLengthCompatLayer — adds Content-Length: 0 for known empty-body API routes
|
// 6. EmptyBodyContentLengthCompatLayer — adds Content-Length: 0 for known empty-body API routes
|
||||||
// 7. CatchPanicLayer — panic → 500
|
// 7. CatchPanicLayer — panic → 500
|
||||||
// 8. ReadinessGateLayer — blocks until ready
|
// 8. RateLimitLayer — conditional (external stack only), per-client 429 throttling
|
||||||
// 9. KeystoneAuthLayer — X-Auth-Token validation
|
// 9. ReadinessGateLayer — blocks until ready
|
||||||
// 10. TraceLayer — request span creation + metrics
|
// 10. KeystoneAuthLayer — X-Auth-Token validation
|
||||||
// 11. RequestLoggingLayer — single completion event per request
|
// 11. TraceLayer — request span creation + metrics
|
||||||
// 12. PropagateRequestIdLayer — X-Request-ID → response
|
// 12. RequestLoggingLayer — single completion event per request
|
||||||
// 13. CompressionLayer — response compression (whitelist, path-aware)
|
// 13. PropagateRequestIdLayer — X-Request-ID → response
|
||||||
// 14. PathCategoryInjectionLayer — injects path category for compression predicate
|
// 14. CompressionLayer — response compression (whitelist, path-aware)
|
||||||
// 15. S3ErrorMessageCompatLayer — missing S3 error message compatibility
|
// 15. PathCategoryInjectionLayer — injects path category for compression predicate
|
||||||
// 16. IcebergRestErrorCompatLayer — Iceberg REST JSON error compatibility
|
// 16. S3ErrorMessageCompatLayer — missing S3 error message compatibility
|
||||||
// 17. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes
|
// 17. IcebergRestErrorCompatLayer — Iceberg REST JSON error compatibility
|
||||||
// 18. ConditionalCorsLayer — S3 API CORS
|
// 18. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes
|
||||||
// 19. RedirectLayer — console redirect (conditional)
|
// 19. ConditionalCorsLayer — S3 API CORS
|
||||||
// 20. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses
|
// 20. RedirectLayer — console redirect (conditional)
|
||||||
// 21. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses
|
// 21. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses
|
||||||
// 22. PublicHealthEndpointLayer — handles public health before s3s host parsing
|
// 22. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses
|
||||||
// 23. VirtualHostStyleHintLayer — actionable error for unroutable virtual-hosted-style (conditional)
|
// 23. PublicHealthEndpointLayer — handles public health before s3s host parsing
|
||||||
// 24. DoubleSlashListBucketsCompatLayer — rewrites `GET //` to `GET /` for ListBuckets (MinIO browser compat)
|
// 24. VirtualHostStyleHintLayer — actionable error for unroutable virtual-hosted-style (conditional)
|
||||||
|
// 25. DoubleSlashListBucketsCompatLayer — rewrites `GET //` to `GET /` for ListBuckets (MinIO browser compat)
|
||||||
// ─────────────────────────────────────────────────────────────
|
// ─────────────────────────────────────────────────────────────
|
||||||
// Batch 1 intentionally keeps the external and internode stacks behaviorally
|
// Batch 1 intentionally keeps the external and internode stacks behaviorally
|
||||||
// identical while giving each path family a named construction boundary.
|
// identical while giving each path family a named construction boundary.
|
||||||
@@ -1223,6 +1259,12 @@ fn process_connection(
|
|||||||
.layer(InternodeRequestContextLiteLayer)
|
.layer(InternodeRequestContextLiteLayer)
|
||||||
.layer(EmptyBodyContentLengthCompatLayer)
|
.layer(EmptyBodyContentLengthCompatLayer)
|
||||||
.layer(CatchPanicLayer::new())
|
.layer(CatchPanicLayer::new())
|
||||||
|
// Per-client API rate limit (backlog#1191): rejects over-limit
|
||||||
|
// requests with 429 before readiness/auth/tracing spend any
|
||||||
|
// work on them, but after the trusted-proxy layer has resolved
|
||||||
|
// a spoof-proof client IP. Absent (None) unless enabled via
|
||||||
|
// RUSTFS_API_RATE_LIMIT_ENABLE with a non-zero RPM.
|
||||||
|
.option_layer(rate_limit_layer.clone())
|
||||||
// CRITICAL: Insert ReadinessGateLayer before business logic
|
// CRITICAL: Insert ReadinessGateLayer before business logic
|
||||||
// This stops requests from hitting IAMAuth or Storage if they are not ready.
|
// This stops requests from hitting IAMAuth or Storage if they are not ready.
|
||||||
.layer(ReadinessGateLayer::new(readiness.clone()))
|
.layer(ReadinessGateLayer::new(readiness.clone()))
|
||||||
|
|||||||
+700
-107
@@ -12,124 +12,391 @@
|
|||||||
// See the License for the specific language governing permissions and
|
// See the License for the specific language governing permissions and
|
||||||
// limitations under the License.
|
// limitations under the License.
|
||||||
|
|
||||||
//! Rate limiting middleware for RustFS
|
//! Per-client request rate limiting for the S3 API (backlog#1191).
|
||||||
//!
|
//!
|
||||||
//! This module provides per-client request rate limiting using a token bucket algorithm.
|
//! A token-bucket limiter keyed by client IP. Over-limit requests are rejected
|
||||||
//! It helps protect against DoS attacks by limiting the number of requests per client.
|
//! with `429 Too Many Requests`, a `Retry-After` header, and the
|
||||||
|
//! `x-ratelimit-*` headers already used by the Swift protocol error mapping.
|
||||||
|
//!
|
||||||
|
//! Scope and invariants:
|
||||||
|
//! - **Opt-in, default off.** [`api_rate_limit_layer_from_env`] returns `None`
|
||||||
|
//! unless `RUSTFS_API_RATE_LIMIT_ENABLE=true` and `RUSTFS_API_RATE_LIMIT_RPM > 0`;
|
||||||
|
//! the layer is then absent from the service stack and the request path is
|
||||||
|
//! unchanged.
|
||||||
|
//! - **Client identity is never taken from request headers.** The key is the
|
||||||
|
//! validated [`ClientInfo::real_ip`] inserted by the trusted-proxy layer
|
||||||
|
//! (which only honors `X-Forwarded-For` from configured proxies) or, absent
|
||||||
|
//! that, the socket peer address ([`RemoteAddr`]). A raw `X-Forwarded-For`
|
||||||
|
//! header can therefore not be used to escape into an attacker-chosen bucket.
|
||||||
|
//! - **Infra traffic is exempt** ([`is_rate_limit_exempt_path`]): health and
|
||||||
|
//! profiling probes, internode RPC/gRPC, and the console (which has its own
|
||||||
|
//! limiter, sharing this module's [`RateLimiter`] core).
|
||||||
|
//! - This is client-facing abuse protection, not internal backpressure — it is
|
||||||
|
//! unrelated to `workload_admission`, which schedules already-admitted work.
|
||||||
|
//!
|
||||||
|
//! Memory is bounded: buckets live in [`SHARD_COUNT`] independently locked
|
||||||
|
//! shards, each capped and periodically swept. A bucket that has been idle
|
||||||
|
//! long enough to have refilled completely is indistinguishable from a fresh
|
||||||
|
//! one, so eviction never makes the limiter more permissive than a restart.
|
||||||
|
|
||||||
use http::{Request, Response};
|
use crate::server::{
|
||||||
|
CONSOLE_PREFIX, FAVICON_PATH, HEALTH_COMPAT_LIVE_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, MINIO_HEALTH_CLUSTER_PATH,
|
||||||
|
MINIO_HEALTH_CLUSTER_READ_PATH, MINIO_HEALTH_LIVE_PATH, MINIO_HEALTH_READY_PATH, PROFILE_CPU_PATH, PROFILE_MEMORY_PATH,
|
||||||
|
RPC_PREFIX, RemoteAddr, TONIC_PREFIX, has_path_prefix,
|
||||||
|
};
|
||||||
|
use bytes::Bytes;
|
||||||
|
use futures::future::{Either, Ready, ready};
|
||||||
|
use http::{HeaderMap, HeaderValue, Request, Response, StatusCode};
|
||||||
|
use http_body_util::{BodyExt, Full};
|
||||||
|
use metrics::counter;
|
||||||
|
use rustfs_trusted_proxies::ClientInfo;
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
|
use std::hash::{BuildHasher, RandomState};
|
||||||
use std::net::IpAddr;
|
use std::net::IpAddr;
|
||||||
use std::sync::{Arc, RwLock};
|
use std::sync::{Arc, Mutex, PoisonError};
|
||||||
use std::time::{Duration, Instant};
|
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||||
use tower::{Layer, Service};
|
use tower::{Layer, Service};
|
||||||
|
use tracing::{debug, warn};
|
||||||
|
|
||||||
/// Configuration for rate limiting
|
const LOG_COMPONENT_SERVER: &str = "server";
|
||||||
#[derive(Debug, Clone)]
|
const LOG_SUBSYSTEM_RATE_LIMIT: &str = "rate_limit";
|
||||||
pub struct RateLimitConfig {
|
const EVENT_REQUEST_RATE_LIMITED: &str = "request_rate_limited";
|
||||||
/// Maximum number of requests per window
|
|
||||||
pub max_requests: u32,
|
pub(crate) const METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL: &str = "rustfs_http_server_requests_rate_limited_total";
|
||||||
/// Time window duration
|
pub(crate) const LABEL_RATE_LIMIT_SCOPE: &str = "scope";
|
||||||
pub window_duration: Duration,
|
pub(crate) const RATE_LIMIT_SCOPE_S3_API: &str = "s3_api";
|
||||||
|
pub(crate) const RATE_LIMIT_SCOPE_CONSOLE: &str = "console";
|
||||||
|
|
||||||
|
// Header names shared with the Swift error mapping (crates/protocols swift/errors.rs).
|
||||||
|
const X_RATE_LIMIT_LIMIT: &str = "x-ratelimit-limit";
|
||||||
|
const X_RATE_LIMIT_REMAINING: &str = "x-ratelimit-remaining";
|
||||||
|
const X_RATE_LIMIT_RESET: &str = "x-ratelimit-reset";
|
||||||
|
const X_REQUEST_ID: &str = "x-request-id";
|
||||||
|
|
||||||
|
/// Shards bound lock contention: each request locks 1/32 of the key space.
|
||||||
|
const SHARD_COUNT: usize = 32;
|
||||||
|
/// Upper bound on tracked client IPs across all shards (~100 bytes each, so
|
||||||
|
/// worst case a few MiB). Real distinct-IP cardinality is bounded by actual
|
||||||
|
/// TCP peers (or validated proxy clients), never by spoofable headers.
|
||||||
|
const MAX_TRACKED_CLIENTS: usize = 100_000;
|
||||||
|
/// Per-shard sweep cadence for dropping refilled-idle buckets.
|
||||||
|
const CLEANUP_INTERVAL: Duration = Duration::from_secs(30);
|
||||||
|
|
||||||
|
/// Sustained-plus-burst request budget for one client.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct RateLimitQuota {
|
||||||
|
/// Sustained budget in requests per minute.
|
||||||
|
pub requests_per_minute: u32,
|
||||||
|
/// Bucket capacity (maximum burst).
|
||||||
|
pub burst: u32,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Default for RateLimitConfig {
|
impl RateLimitQuota {
|
||||||
fn default() -> Self {
|
/// Build a quota from operator configuration. `rpm == 0` means unlimited
|
||||||
Self {
|
/// (`None`); `burst == 0` means "same as RPM".
|
||||||
max_requests: 100,
|
pub fn per_minute(rpm: u32, burst: u32) -> Option<Self> {
|
||||||
window_duration: Duration::from_secs(60),
|
if rpm == 0 {
|
||||||
|
return None;
|
||||||
}
|
}
|
||||||
|
Some(Self {
|
||||||
|
requests_per_minute: rpm,
|
||||||
|
burst: if burst == 0 { rpm } else { burst },
|
||||||
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Token bucket for rate limiting
|
/// Delay hints for a rejected request, in whole seconds (both at least 1).
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub struct ThrottleInfo {
|
||||||
|
/// Seconds until at least one token is available again (`Retry-After`).
|
||||||
|
pub retry_after_secs: u64,
|
||||||
|
/// Seconds until the bucket is completely full (`x-ratelimit-reset`).
|
||||||
|
pub reset_after_secs: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
pub enum RateLimitDecision {
|
||||||
|
Allowed,
|
||||||
|
Limited(ThrottleInfo),
|
||||||
|
}
|
||||||
|
|
||||||
|
/// One client's bucket. Capacity and refill rate are limiter-wide (single
|
||||||
|
/// quota for all clients), so they are not duplicated per entry.
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
struct TokenBucket {
|
struct TokenBucket {
|
||||||
/// Maximum number of tokens
|
|
||||||
max_tokens: u32,
|
|
||||||
/// Current number of tokens
|
|
||||||
tokens: f64,
|
tokens: f64,
|
||||||
/// Rate at which tokens are refilled (tokens per second)
|
|
||||||
refill_rate: f64,
|
|
||||||
/// Last time tokens were refilled
|
|
||||||
last_refill: Instant,
|
last_refill: Instant,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl TokenBucket {
|
impl TokenBucket {
|
||||||
fn new(max_tokens: u32, window_duration: Duration) -> Self {
|
fn full(capacity: f64, now: Instant) -> Self {
|
||||||
Self {
|
Self {
|
||||||
max_tokens,
|
tokens: capacity,
|
||||||
tokens: max_tokens as f64,
|
last_refill: now,
|
||||||
refill_rate: max_tokens as f64 / window_duration.as_secs_f64(),
|
|
||||||
last_refill: Instant::now(),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn try_consume(&mut self) -> bool {
|
fn try_consume(&mut self, capacity: f64, refill_per_sec: f64, now: Instant) -> RateLimitDecision {
|
||||||
self.refill();
|
let elapsed = now.saturating_duration_since(self.last_refill).as_secs_f64();
|
||||||
|
self.tokens = (self.tokens + elapsed * refill_per_sec).min(capacity);
|
||||||
|
self.last_refill = now;
|
||||||
|
|
||||||
if self.tokens >= 1.0 {
|
if self.tokens >= 1.0 {
|
||||||
self.tokens -= 1.0;
|
self.tokens -= 1.0;
|
||||||
true
|
RateLimitDecision::Allowed
|
||||||
} else {
|
} else {
|
||||||
false
|
let retry_after_secs = (((1.0 - self.tokens) / refill_per_sec).ceil() as u64).max(1);
|
||||||
|
let reset_after_secs = (((capacity - self.tokens) / refill_per_sec).ceil() as u64).max(1);
|
||||||
|
RateLimitDecision::Limited(ThrottleInfo {
|
||||||
|
retry_after_secs,
|
||||||
|
reset_after_secs,
|
||||||
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn refill(&mut self) {
|
/// Idle long enough that refill would have restored a full bucket: the
|
||||||
let now = Instant::now();
|
/// entry carries no information a fresh one would not, so it can be
|
||||||
let elapsed = now.duration_since(self.last_refill).as_secs_f64();
|
/// dropped losslessly.
|
||||||
self.tokens = (self.tokens + elapsed * self.refill_rate).min(self.max_tokens as f64);
|
fn is_refilled_idle(&self, capacity: f64, refill_per_sec: f64, now: Instant) -> bool {
|
||||||
self.last_refill = now;
|
now.saturating_duration_since(self.last_refill).as_secs_f64() * refill_per_sec >= capacity
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Rate limiter state
|
|
||||||
#[derive(Debug)]
|
#[derive(Debug)]
|
||||||
struct RateLimiterState {
|
struct Shard {
|
||||||
/// Per-client token buckets
|
|
||||||
clients: HashMap<IpAddr, TokenBucket>,
|
clients: HashMap<IpAddr, TokenBucket>,
|
||||||
/// Configuration
|
next_cleanup: Instant,
|
||||||
config: RateLimitConfig,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
impl RateLimiterState {
|
/// Sharded per-IP token-bucket limiter. Shared by the S3 API tower layer and
|
||||||
fn new(config: RateLimitConfig) -> Self {
|
/// the console middleware (each with its own instance and quota).
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct RateLimiter {
|
||||||
|
shards: Vec<Mutex<Shard>>,
|
||||||
|
shard_hasher: RandomState,
|
||||||
|
quota: RateLimitQuota,
|
||||||
|
capacity: f64,
|
||||||
|
refill_per_sec: f64,
|
||||||
|
max_clients_per_shard: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl RateLimiter {
|
||||||
|
pub fn new(quota: RateLimitQuota) -> Self {
|
||||||
|
Self::with_limits(quota, SHARD_COUNT, MAX_TRACKED_CLIENTS)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn with_limits(quota: RateLimitQuota, shard_count: usize, max_clients: usize) -> Self {
|
||||||
|
let now = Instant::now();
|
||||||
Self {
|
Self {
|
||||||
clients: HashMap::new(),
|
shards: (0..shard_count.max(1))
|
||||||
config,
|
.map(|_| {
|
||||||
|
Mutex::new(Shard {
|
||||||
|
clients: HashMap::new(),
|
||||||
|
next_cleanup: now + CLEANUP_INTERVAL,
|
||||||
|
})
|
||||||
|
})
|
||||||
|
.collect(),
|
||||||
|
shard_hasher: RandomState::new(),
|
||||||
|
quota,
|
||||||
|
capacity: f64::from(quota.burst.max(1)),
|
||||||
|
refill_per_sec: f64::from(quota.requests_per_minute) / 60.0,
|
||||||
|
max_clients_per_shard: (max_clients / shard_count.max(1)).max(1),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn check_rate_limit(&mut self, client_ip: IpAddr) -> bool {
|
pub fn quota(&self) -> RateLimitQuota {
|
||||||
let bucket = self
|
self.quota
|
||||||
.clients
|
|
||||||
.entry(client_ip)
|
|
||||||
.or_insert_with(|| TokenBucket::new(self.config.max_requests, self.config.window_duration));
|
|
||||||
|
|
||||||
bucket.try_consume()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fn cleanup_expired(&mut self) {
|
/// Consume one token for `ip`, creating its bucket on first sight.
|
||||||
let now = Instant::now();
|
pub fn check(&self, ip: IpAddr) -> RateLimitDecision {
|
||||||
let window = self.config.window_duration;
|
self.check_at(ip, Instant::now())
|
||||||
self.clients
|
}
|
||||||
.retain(|_, bucket| now.duration_since(bucket.last_refill) < window * 2);
|
|
||||||
|
fn check_at(&self, ip: IpAddr, now: Instant) -> RateLimitDecision {
|
||||||
|
let shard = &self.shards[self.shard_index(ip)];
|
||||||
|
let mut shard = shard.lock().unwrap_or_else(PoisonError::into_inner);
|
||||||
|
|
||||||
|
if now >= shard.next_cleanup {
|
||||||
|
let (capacity, refill) = (self.capacity, self.refill_per_sec);
|
||||||
|
shard
|
||||||
|
.clients
|
||||||
|
.retain(|_, bucket| !bucket.is_refilled_idle(capacity, refill, now));
|
||||||
|
shard.next_cleanup = now + CLEANUP_INTERVAL;
|
||||||
|
}
|
||||||
|
|
||||||
|
if let Some(bucket) = shard.clients.get_mut(&ip) {
|
||||||
|
return bucket.try_consume(self.capacity, self.refill_per_sec, now);
|
||||||
|
}
|
||||||
|
|
||||||
|
if shard.clients.len() >= self.max_clients_per_shard {
|
||||||
|
// At capacity: reclaim the most idle entry (closest to full, so the
|
||||||
|
// least information is lost). O(shard len) but only under a
|
||||||
|
// distinct-IP flood, where correctness beats micro-latency.
|
||||||
|
if let Some(stalest) = shard
|
||||||
|
.clients
|
||||||
|
.iter()
|
||||||
|
.min_by_key(|(_, bucket)| bucket.last_refill)
|
||||||
|
.map(|(ip, _)| *ip)
|
||||||
|
{
|
||||||
|
shard.clients.remove(&stalest);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
shard
|
||||||
|
.clients
|
||||||
|
.entry(ip)
|
||||||
|
.or_insert_with(|| TokenBucket::full(self.capacity, now))
|
||||||
|
.try_consume(self.capacity, self.refill_per_sec, now)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn shard_index(&self, ip: IpAddr) -> usize {
|
||||||
|
(self.shard_hasher.hash_one(ip) as usize) % self.shards.len()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn tracked_clients(&self) -> usize {
|
||||||
|
self.shards
|
||||||
|
.iter()
|
||||||
|
.map(|shard| shard.lock().unwrap_or_else(PoisonError::into_inner).clients.len())
|
||||||
|
.sum()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Rate limiting layer
|
/// Resolve the client identity for rate limiting.
|
||||||
|
///
|
||||||
|
/// Only trusted sources are consulted: `ClientInfo` is inserted by the
|
||||||
|
/// trusted-proxy layer after validating the proxy chain, and `RemoteAddr` is
|
||||||
|
/// the raw socket peer. Spoofable headers (`X-Forwarded-For`, `X-Real-IP`)
|
||||||
|
/// are deliberately never read here.
|
||||||
|
pub(crate) fn client_ip<B>(req: &Request<B>) -> Option<IpAddr> {
|
||||||
|
if let Some(info) = req.extensions().get::<ClientInfo>() {
|
||||||
|
return Some(info.real_ip);
|
||||||
|
}
|
||||||
|
req.extensions().get::<RemoteAddr>().map(|remote| remote.0.ip())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Paths exempt from S3 API rate limiting.
|
||||||
|
///
|
||||||
|
/// - Health/profiling probes: kubelet and load-balancer checks must never be
|
||||||
|
/// throttled, or the limiter itself becomes an availability risk.
|
||||||
|
/// - Internode RPC/gRPC: cluster-internal traffic; throttling it would let an
|
||||||
|
/// external flood starve replication and heal.
|
||||||
|
/// - Console: has its own dedicated limiter (`RUSTFS_CONSOLE_RATE_LIMIT_*`).
|
||||||
|
fn is_rate_limit_exempt_path(path: &str) -> bool {
|
||||||
|
matches!(
|
||||||
|
path,
|
||||||
|
HEALTH_PREFIX
|
||||||
|
| HEALTH_COMPAT_LIVE_PATH
|
||||||
|
| HEALTH_READY_PATH
|
||||||
|
| MINIO_HEALTH_LIVE_PATH
|
||||||
|
| MINIO_HEALTH_READY_PATH
|
||||||
|
| MINIO_HEALTH_CLUSTER_PATH
|
||||||
|
| MINIO_HEALTH_CLUSTER_READ_PATH
|
||||||
|
| PROFILE_CPU_PATH
|
||||||
|
| PROFILE_MEMORY_PATH
|
||||||
|
| FAVICON_PATH
|
||||||
|
) || has_path_prefix(path, RPC_PREFIX)
|
||||||
|
|| has_path_prefix(path, TONIC_PREFIX)
|
||||||
|
|| has_path_prefix(path, CONSOLE_PREFIX)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Apply the standard throttling headers shared by every rate-limited scope.
|
||||||
|
pub(crate) fn apply_throttle_headers(headers: &mut HeaderMap, limit_rpm: u32, throttle: &ThrottleInfo) {
|
||||||
|
let reset_unix = SystemTime::now()
|
||||||
|
.duration_since(UNIX_EPOCH)
|
||||||
|
.unwrap_or_default()
|
||||||
|
.as_secs()
|
||||||
|
.saturating_add(throttle.reset_after_secs);
|
||||||
|
headers.insert(http::header::RETRY_AFTER, u64_header(throttle.retry_after_secs));
|
||||||
|
headers.insert(X_RATE_LIMIT_LIMIT, u64_header(u64::from(limit_rpm)));
|
||||||
|
headers.insert(X_RATE_LIMIT_REMAINING, HeaderValue::from_static("0"));
|
||||||
|
headers.insert(X_RATE_LIMIT_RESET, u64_header(reset_unix));
|
||||||
|
}
|
||||||
|
|
||||||
|
fn u64_header(value: u64) -> HeaderValue {
|
||||||
|
HeaderValue::from_str(&value.to_string()).unwrap_or_else(|_| HeaderValue::from_static("0"))
|
||||||
|
}
|
||||||
|
|
||||||
|
type BoxError = Box<dyn std::error::Error + Send + Sync>;
|
||||||
|
type BoxBody = http_body_util::combinators::UnsyncBoxBody<Bytes, BoxError>;
|
||||||
|
|
||||||
|
/// Build the S3-style `429` rejection. The request id (already generated by
|
||||||
|
/// `SetRequestIdLayer`, which sits outside this layer) is echoed manually
|
||||||
|
/// because the rejection short-circuits below `PropagateRequestIdLayer`.
|
||||||
|
fn s3_too_many_requests_response(request_id: Option<&HeaderValue>, limit_rpm: u32, throttle: &ThrottleInfo) -> Response<BoxBody> {
|
||||||
|
// The header may be client-supplied (SetRequestIdLayer only fills it when
|
||||||
|
// absent), so gate the XML interpolation on a UUID-safe charset instead of
|
||||||
|
// reflecting arbitrary bytes into the body.
|
||||||
|
let request_id_xml = request_id
|
||||||
|
.and_then(|value| value.to_str().ok())
|
||||||
|
.filter(|id| !id.is_empty() && id.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'-'))
|
||||||
|
.map(|id| format!("<RequestId>{id}</RequestId>"))
|
||||||
|
.unwrap_or_default();
|
||||||
|
let body = format!(
|
||||||
|
"<?xml version=\"1.0\" encoding=\"UTF-8\"?>\
|
||||||
|
<Error><Code>TooManyRequests</Code>\
|
||||||
|
<Message>Request rate limit exceeded. Reduce your request rate and retry after the indicated delay.</Message>\
|
||||||
|
{request_id_xml}</Error>"
|
||||||
|
);
|
||||||
|
let body: BoxBody = Full::new(Bytes::from(body))
|
||||||
|
.map_err(|e| -> BoxError { Box::new(e) })
|
||||||
|
.boxed_unsync();
|
||||||
|
|
||||||
|
let mut response = Response::new(body);
|
||||||
|
*response.status_mut() = StatusCode::TOO_MANY_REQUESTS;
|
||||||
|
response
|
||||||
|
.headers_mut()
|
||||||
|
.insert(http::header::CONTENT_TYPE, HeaderValue::from_static("application/xml"));
|
||||||
|
apply_throttle_headers(response.headers_mut(), limit_rpm, throttle);
|
||||||
|
if let Some(id) = request_id {
|
||||||
|
response.headers_mut().insert(X_REQUEST_ID, id.clone());
|
||||||
|
}
|
||||||
|
response
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Build the S3 API rate limit layer from `RUSTFS_API_RATE_LIMIT_*`.
|
||||||
|
///
|
||||||
|
/// Returns `None` (no layer in the stack, zero request-path change) unless
|
||||||
|
/// explicitly enabled with a non-zero RPM.
|
||||||
|
pub fn api_rate_limit_layer_from_env() -> Option<RateLimitLayer> {
|
||||||
|
if !rustfs_utils::get_env_bool(rustfs_config::ENV_API_RATE_LIMIT_ENABLE, rustfs_config::DEFAULT_API_RATE_LIMIT_ENABLE) {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let rpm = rustfs_utils::get_env_u32(rustfs_config::ENV_API_RATE_LIMIT_RPM, rustfs_config::DEFAULT_API_RATE_LIMIT_RPM);
|
||||||
|
let burst = rustfs_utils::get_env_u32(rustfs_config::ENV_API_RATE_LIMIT_BURST, rustfs_config::DEFAULT_API_RATE_LIMIT_BURST);
|
||||||
|
let Some(quota) = RateLimitQuota::per_minute(rpm, burst) else {
|
||||||
|
warn!(
|
||||||
|
component = LOG_COMPONENT_SERVER,
|
||||||
|
subsystem = LOG_SUBSYSTEM_RATE_LIMIT,
|
||||||
|
"{} is enabled but {} is 0; API rate limiting stays inactive",
|
||||||
|
rustfs_config::ENV_API_RATE_LIMIT_ENABLE,
|
||||||
|
rustfs_config::ENV_API_RATE_LIMIT_RPM
|
||||||
|
);
|
||||||
|
return None;
|
||||||
|
};
|
||||||
|
Some(RateLimitLayer::new(quota))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Tower layer enforcing the per-IP quota on the external S3 service stack.
|
||||||
|
///
|
||||||
|
/// All clones share one [`RateLimiter`], so per-connection stack construction
|
||||||
|
/// keeps a single global view of client budgets.
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct RateLimitLayer {
|
pub struct RateLimitLayer {
|
||||||
state: Arc<RwLock<RateLimiterState>>,
|
limiter: Arc<RateLimiter>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl RateLimitLayer {
|
impl RateLimitLayer {
|
||||||
pub fn new(config: RateLimitConfig) -> Self {
|
pub fn new(quota: RateLimitQuota) -> Self {
|
||||||
Self {
|
Self {
|
||||||
state: Arc::new(RwLock::new(RateLimiterState::new(config))),
|
limiter: Arc::new(RateLimiter::new(quota)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn quota(&self) -> RateLimitQuota {
|
||||||
|
self.limiter.quota()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S> Layer<S> for RateLimitLayer {
|
impl<S> Layer<S> for RateLimitLayer {
|
||||||
@@ -138,75 +405,401 @@ impl<S> Layer<S> for RateLimitLayer {
|
|||||||
fn layer(&self, inner: S) -> Self::Service {
|
fn layer(&self, inner: S) -> Self::Service {
|
||||||
RateLimitService {
|
RateLimitService {
|
||||||
inner,
|
inner,
|
||||||
state: self.state.clone(),
|
limiter: self.limiter.clone(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Rate limiting service
|
/// See [`RateLimitLayer`]. The response type is pinned to the stack's
|
||||||
|
/// `BoxBody` so the layer composes with `option_layer` at its position just
|
||||||
|
/// outside `ReadinessGateLayer` (which already boxes response bodies).
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct RateLimitService<S> {
|
pub struct RateLimitService<S> {
|
||||||
inner: S,
|
inner: S,
|
||||||
state: Arc<RwLock<RateLimiterState>>,
|
limiter: Arc<RateLimiter>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<S, ReqBody, ResBody> Service<Request<ReqBody>> for RateLimitService<S>
|
impl<S, ReqBody> Service<Request<ReqBody>> for RateLimitService<S>
|
||||||
where
|
where
|
||||||
S: Service<Request<ReqBody>, Response = Response<ResBody>>,
|
S: Service<Request<ReqBody>, Response = Response<BoxBody>>,
|
||||||
{
|
{
|
||||||
type Response = S::Response;
|
type Response = Response<BoxBody>;
|
||||||
type Error = S::Error;
|
type Error = S::Error;
|
||||||
type Future = S::Future;
|
type Future = Either<S::Future, Ready<Result<Response<BoxBody>, S::Error>>>;
|
||||||
|
|
||||||
fn poll_ready(&mut self, cx: &mut std::task::Context<'_>) -> std::task::Poll<Result<(), Self::Error>> {
|
fn poll_ready(&mut self, cx: &mut std::task::Context<'_>) -> std::task::Poll<Result<(), Self::Error>> {
|
||||||
self.inner.poll_ready(cx)
|
self.inner.poll_ready(cx)
|
||||||
}
|
}
|
||||||
|
|
||||||
fn call(&mut self, req: Request<ReqBody>) -> Self::Future {
|
fn call(&mut self, req: Request<ReqBody>) -> Self::Future {
|
||||||
// Extract client IP from request
|
if is_rate_limit_exempt_path(req.uri().path()) {
|
||||||
let client_ip = extract_client_ip(&req);
|
return Either::Left(self.inner.call(req));
|
||||||
|
}
|
||||||
// Check rate limit
|
let Some(ip) = client_ip(&req) else {
|
||||||
let allowed = {
|
// Fail open: without a socket peer address there is no trustworthy
|
||||||
let mut state = self.state.write().unwrap_or_else(|e| e.into_inner());
|
// client identity, and deriving one from spoofable headers would
|
||||||
// Periodically cleanup expired entries
|
// let clients pick their own bucket.
|
||||||
if rand::random::<u32>().is_multiple_of(100) {
|
return Either::Left(self.inner.call(req));
|
||||||
state.cleanup_expired();
|
|
||||||
}
|
|
||||||
state.check_rate_limit(client_ip)
|
|
||||||
};
|
};
|
||||||
|
|
||||||
if allowed {
|
match self.limiter.check(ip) {
|
||||||
self.inner.call(req)
|
RateLimitDecision::Allowed => Either::Left(self.inner.call(req)),
|
||||||
} else {
|
RateLimitDecision::Limited(throttle) => {
|
||||||
// Return 429 Too Many Requests
|
counter!(
|
||||||
// Note: In a real implementation, you'd want to return a proper 429 response.
|
METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL,
|
||||||
// For now, we just pass the request through to the inner service.
|
LABEL_RATE_LIMIT_SCOPE => RATE_LIMIT_SCOPE_S3_API
|
||||||
// The rate limiting is enforced at the application level.
|
)
|
||||||
self.inner.call(req)
|
.increment(1);
|
||||||
|
debug!(
|
||||||
|
event = EVENT_REQUEST_RATE_LIMITED,
|
||||||
|
component = LOG_COMPONENT_SERVER,
|
||||||
|
subsystem = LOG_SUBSYSTEM_RATE_LIMIT,
|
||||||
|
scope = RATE_LIMIT_SCOPE_S3_API,
|
||||||
|
client_ip = %ip,
|
||||||
|
retry_after_secs = throttle.retry_after_secs,
|
||||||
|
"Request rejected by API rate limit"
|
||||||
|
);
|
||||||
|
let limit_rpm = self.limiter.quota().requests_per_minute;
|
||||||
|
Either::Right(ready(Ok(s3_too_many_requests_response(
|
||||||
|
req.headers().get(X_REQUEST_ID),
|
||||||
|
limit_rpm,
|
||||||
|
&throttle,
|
||||||
|
))))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Extract client IP from request headers or connection info
|
#[cfg(test)]
|
||||||
fn extract_client_ip<ReqBody>(req: &Request<ReqBody>) -> IpAddr {
|
mod tests {
|
||||||
// Try X-Forwarded-For header first
|
use super::*;
|
||||||
if let Some(forwarded) = req.headers().get("x-forwarded-for")
|
use http_body_util::BodyExt;
|
||||||
&& let Ok(forwarded_str) = forwarded.to_str()
|
use serial_test::serial;
|
||||||
&& let Some(first_ip) = forwarded_str.split(',').next()
|
use std::net::{Ipv4Addr, SocketAddr};
|
||||||
&& let Ok(ip) = first_ip.trim().parse::<IpAddr>()
|
use std::task::{Context, Poll};
|
||||||
{
|
|
||||||
return ip;
|
fn quota(rpm: u32, burst: u32) -> RateLimitQuota {
|
||||||
|
RateLimitQuota::per_minute(rpm, burst).expect("non-zero rpm")
|
||||||
}
|
}
|
||||||
|
|
||||||
// Try X-Real-IP header
|
fn ip(last: u8) -> IpAddr {
|
||||||
if let Some(real_ip) = req.headers().get("x-real-ip")
|
IpAddr::V4(Ipv4Addr::new(203, 0, 113, last))
|
||||||
&& let Ok(ip_str) = real_ip.to_str()
|
|
||||||
&& let Ok(ip) = ip_str.parse::<IpAddr>()
|
|
||||||
{
|
|
||||||
return ip;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// Default to localhost
|
#[test]
|
||||||
IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)
|
fn quota_zero_rpm_means_unlimited_and_zero_burst_defaults_to_rpm() {
|
||||||
|
assert_eq!(RateLimitQuota::per_minute(0, 50), None);
|
||||||
|
assert_eq!(quota(120, 0).burst, 120);
|
||||||
|
assert_eq!(quota(120, 10).burst, 10);
|
||||||
|
assert_eq!(quota(120, 500).burst, 500);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn bucket_exhaustion_returns_429_hints_and_window_recovers() {
|
||||||
|
let limiter = RateLimiter::new(quota(60, 5)); // 1 token/sec, burst 5
|
||||||
|
let start = Instant::now();
|
||||||
|
|
||||||
|
for _ in 0..5 {
|
||||||
|
assert_eq!(limiter.check_at(ip(1), start), RateLimitDecision::Allowed);
|
||||||
|
}
|
||||||
|
let RateLimitDecision::Limited(throttle) = limiter.check_at(ip(1), start) else {
|
||||||
|
panic!("sixth request within the same instant must be limited");
|
||||||
|
};
|
||||||
|
assert_eq!(throttle.retry_after_secs, 1);
|
||||||
|
assert_eq!(throttle.reset_after_secs, 5);
|
||||||
|
|
||||||
|
// One refill interval restores exactly one token.
|
||||||
|
assert_eq!(limiter.check_at(ip(1), start + Duration::from_secs(1)), RateLimitDecision::Allowed);
|
||||||
|
assert!(matches!(
|
||||||
|
limiter.check_at(ip(1), start + Duration::from_secs(1)),
|
||||||
|
RateLimitDecision::Limited(_)
|
||||||
|
));
|
||||||
|
|
||||||
|
// A full window restores the full burst.
|
||||||
|
let later = start + Duration::from_secs(61);
|
||||||
|
for _ in 0..5 {
|
||||||
|
assert_eq!(limiter.check_at(ip(1), later), RateLimitDecision::Allowed);
|
||||||
|
}
|
||||||
|
assert!(matches!(limiter.check_at(ip(1), later), RateLimitDecision::Limited(_)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn clients_have_independent_buckets() {
|
||||||
|
let limiter = RateLimiter::new(quota(60, 1));
|
||||||
|
let now = Instant::now();
|
||||||
|
assert_eq!(limiter.check_at(ip(1), now), RateLimitDecision::Allowed);
|
||||||
|
assert!(matches!(limiter.check_at(ip(1), now), RateLimitDecision::Limited(_)));
|
||||||
|
assert_eq!(limiter.check_at(ip(2), now), RateLimitDecision::Allowed);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn tracked_clients_stay_bounded_under_distinct_ip_flood() {
|
||||||
|
let limiter = RateLimiter::with_limits(quota(60, 1), 4, 64);
|
||||||
|
let now = Instant::now();
|
||||||
|
for i in 0..10_000u32 {
|
||||||
|
limiter.check_at(IpAddr::V4(Ipv4Addr::from(i)), now);
|
||||||
|
}
|
||||||
|
assert!(
|
||||||
|
limiter.tracked_clients() <= 64,
|
||||||
|
"tracked {} clients, cap is 64",
|
||||||
|
limiter.tracked_clients()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn refilled_idle_buckets_are_swept() {
|
||||||
|
let limiter = RateLimiter::with_limits(quota(60, 5), 1, 100);
|
||||||
|
let start = Instant::now();
|
||||||
|
limiter.check_at(ip(1), start);
|
||||||
|
assert_eq!(limiter.tracked_clients(), 1);
|
||||||
|
|
||||||
|
// After burst/rate = 5s the bucket is full again; the next request past
|
||||||
|
// the cleanup deadline sweeps it.
|
||||||
|
let after_sweep = start + CLEANUP_INTERVAL + Duration::from_secs(6);
|
||||||
|
limiter.check_at(ip(2), after_sweep);
|
||||||
|
assert_eq!(limiter.tracked_clients(), 1, "idle refilled bucket must be dropped");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn concurrent_hammering_never_exceeds_burst() {
|
||||||
|
let limiter = Arc::new(RateLimiter::new(quota(60, 100)));
|
||||||
|
let now = Instant::now();
|
||||||
|
let allowed = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||||
|
|
||||||
|
std::thread::scope(|scope| {
|
||||||
|
for _ in 0..8 {
|
||||||
|
let limiter = Arc::clone(&limiter);
|
||||||
|
let allowed = Arc::clone(&allowed);
|
||||||
|
scope.spawn(move || {
|
||||||
|
for _ in 0..50 {
|
||||||
|
if limiter.check_at(ip(9), now) == RateLimitDecision::Allowed {
|
||||||
|
allowed.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
});
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
// 8 threads × 50 attempts = 400 attempts at the same instant: exactly
|
||||||
|
// the burst may pass, regardless of interleaving.
|
||||||
|
assert_eq!(allowed.load(std::sync::atomic::Ordering::Relaxed), 100);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn exempt_paths_cover_probes_internode_and_console_only() {
|
||||||
|
for path in [
|
||||||
|
"/health",
|
||||||
|
"/health/ready",
|
||||||
|
"/minio/health/live",
|
||||||
|
"/minio/health/cluster/read",
|
||||||
|
"/profile/cpu",
|
||||||
|
"/favicon.ico",
|
||||||
|
"/rustfs/rpc/anything",
|
||||||
|
"/node_service.NodeService/Ping",
|
||||||
|
"/rustfs/console/index.html",
|
||||||
|
] {
|
||||||
|
assert!(is_rate_limit_exempt_path(path), "{path} must be exempt");
|
||||||
|
}
|
||||||
|
for path in [
|
||||||
|
"/bucket/object",
|
||||||
|
"/",
|
||||||
|
"/rustfs/admin/v3/info",
|
||||||
|
"/minio/admin/v3/info",
|
||||||
|
"/healthy-bucket/object", // prefix boundary: not /health
|
||||||
|
"/iceberg/v1/config",
|
||||||
|
] {
|
||||||
|
assert!(!is_rate_limit_exempt_path(path), "{path} must be limited");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ---- service-level tests ----
|
||||||
|
|
||||||
|
#[derive(Clone)]
|
||||||
|
struct OkService;
|
||||||
|
|
||||||
|
impl<ReqBody> Service<Request<ReqBody>> for OkService {
|
||||||
|
type Response = Response<BoxBody>;
|
||||||
|
type Error = std::convert::Infallible;
|
||||||
|
type Future = Ready<Result<Self::Response, Self::Error>>;
|
||||||
|
|
||||||
|
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
|
||||||
|
Poll::Ready(Ok(()))
|
||||||
|
}
|
||||||
|
|
||||||
|
fn call(&mut self, _req: Request<ReqBody>) -> Self::Future {
|
||||||
|
let body: BoxBody = Full::new(Bytes::from_static(b"ok"))
|
||||||
|
.map_err(|e| -> BoxError { Box::new(e) })
|
||||||
|
.boxed_unsync();
|
||||||
|
ready(Ok(Response::new(body)))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn service_with_quota(rpm: u32, burst: u32) -> RateLimitService<OkService> {
|
||||||
|
RateLimitLayer::new(quota(rpm, burst)).layer(OkService)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn request_from(peer: IpAddr, path: &str) -> Request<()> {
|
||||||
|
let mut req = Request::builder().uri(path).body(()).expect("request");
|
||||||
|
req.extensions_mut().insert(RemoteAddr(SocketAddr::new(peer, 51000)));
|
||||||
|
req
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn over_limit_returns_429_with_retry_after_and_ratelimit_headers() {
|
||||||
|
let mut service = service_with_quota(60, 2);
|
||||||
|
|
||||||
|
for _ in 0..2 {
|
||||||
|
let resp = service.call(request_from(ip(1), "/bucket/object")).await.expect("ok");
|
||||||
|
assert_eq!(resp.status(), StatusCode::OK);
|
||||||
|
}
|
||||||
|
|
||||||
|
let mut req = request_from(ip(1), "/bucket/object");
|
||||||
|
req.headers_mut().insert(X_REQUEST_ID, HeaderValue::from_static("req-123"));
|
||||||
|
let resp = service.call(req).await.expect("ok");
|
||||||
|
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
assert_eq!(resp.headers().get(http::header::RETRY_AFTER).and_then(|v| v.to_str().ok()), Some("1"));
|
||||||
|
assert_eq!(resp.headers().get(X_RATE_LIMIT_LIMIT).and_then(|v| v.to_str().ok()), Some("60"));
|
||||||
|
assert_eq!(resp.headers().get(X_RATE_LIMIT_REMAINING).and_then(|v| v.to_str().ok()), Some("0"));
|
||||||
|
assert!(resp.headers().contains_key(X_RATE_LIMIT_RESET));
|
||||||
|
assert_eq!(resp.headers().get(X_REQUEST_ID).and_then(|v| v.to_str().ok()), Some("req-123"));
|
||||||
|
assert_eq!(
|
||||||
|
resp.headers().get(http::header::CONTENT_TYPE).and_then(|v| v.to_str().ok()),
|
||||||
|
Some("application/xml")
|
||||||
|
);
|
||||||
|
|
||||||
|
let body = resp.into_body().collect().await.expect("body").to_bytes();
|
||||||
|
let body = String::from_utf8_lossy(&body);
|
||||||
|
assert!(body.contains("<Code>TooManyRequests</Code>"), "body: {body}");
|
||||||
|
assert!(body.contains("<RequestId>req-123</RequestId>"), "body: {body}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn hostile_request_id_is_not_reflected_into_the_429_body() {
|
||||||
|
let mut service = service_with_quota(60, 1);
|
||||||
|
let _ = service.call(request_from(ip(3), "/bucket/object")).await.expect("ok");
|
||||||
|
|
||||||
|
let mut req = request_from(ip(3), "/bucket/object");
|
||||||
|
req.headers_mut()
|
||||||
|
.insert(X_REQUEST_ID, HeaderValue::from_static("<Code>evil</Code>"));
|
||||||
|
let resp = service.call(req).await.expect("ok");
|
||||||
|
assert_eq!(resp.status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
|
||||||
|
let body = resp.into_body().collect().await.expect("body").to_bytes();
|
||||||
|
let body = String::from_utf8_lossy(&body);
|
||||||
|
assert!(!body.contains("evil"), "client-controlled request id must not be reflected: {body}");
|
||||||
|
assert!(!body.contains("<RequestId>"), "malformed id must be omitted entirely: {body}");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn exempt_paths_are_never_limited() {
|
||||||
|
let mut service = service_with_quota(60, 1);
|
||||||
|
for _ in 0..5 {
|
||||||
|
let resp = service.call(request_from(ip(1), "/health/ready")).await.expect("ok");
|
||||||
|
assert_eq!(resp.status(), StatusCode::OK);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn forwarded_for_header_cannot_choose_the_bucket() {
|
||||||
|
let mut service = service_with_quota(60, 1);
|
||||||
|
|
||||||
|
// Same socket peer, rotating spoofed headers: still one bucket.
|
||||||
|
let mut first = request_from(ip(1), "/bucket/object");
|
||||||
|
first
|
||||||
|
.headers_mut()
|
||||||
|
.insert("x-forwarded-for", HeaderValue::from_static("10.0.0.1"));
|
||||||
|
assert_eq!(service.call(first).await.expect("ok").status(), StatusCode::OK);
|
||||||
|
|
||||||
|
let mut second = request_from(ip(1), "/bucket/object");
|
||||||
|
second
|
||||||
|
.headers_mut()
|
||||||
|
.insert("x-forwarded-for", HeaderValue::from_static("10.0.0.2"));
|
||||||
|
second.headers_mut().insert("x-real-ip", HeaderValue::from_static("10.0.0.3"));
|
||||||
|
assert_eq!(service.call(second).await.expect("ok").status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn validated_client_info_takes_precedence_over_socket_peer() {
|
||||||
|
let mut service = service_with_quota(60, 1);
|
||||||
|
|
||||||
|
// Two clients behind the same trusted proxy socket: distinct buckets.
|
||||||
|
for client in [ip(10), ip(11)] {
|
||||||
|
let mut req = request_from(ip(1), "/bucket/object");
|
||||||
|
req.extensions_mut().insert(ClientInfo::direct(SocketAddr::new(client, 443)));
|
||||||
|
assert_eq!(service.call(req).await.expect("ok").status(), StatusCode::OK);
|
||||||
|
}
|
||||||
|
|
||||||
|
// And the validated identity is throttled independently of the socket.
|
||||||
|
let mut req = request_from(ip(2), "/bucket/object");
|
||||||
|
req.extensions_mut().insert(ClientInfo::direct(SocketAddr::new(ip(10), 443)));
|
||||||
|
assert_eq!(service.call(req).await.expect("ok").status(), StatusCode::TOO_MANY_REQUESTS);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn missing_client_identity_fails_open() {
|
||||||
|
let mut service = service_with_quota(60, 1);
|
||||||
|
for _ in 0..3 {
|
||||||
|
let req = Request::builder().uri("/bucket/object").body(()).expect("request");
|
||||||
|
assert_eq!(service.call(req).await.expect("ok").status(), StatusCode::OK);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
#[serial]
|
||||||
|
fn env_constructor_defaults_to_disabled() {
|
||||||
|
temp_env::with_vars(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_ENABLE, None::<&str>),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_RPM, None),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_BURST, None),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
assert!(api_rate_limit_layer_from_env().is_none());
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
#[serial]
|
||||||
|
fn env_constructor_requires_enable_and_nonzero_rpm() {
|
||||||
|
temp_env::with_vars(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_ENABLE, Some("true")),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_RPM, None),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_BURST, None),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
assert!(api_rate_limit_layer_from_env().is_none(), "enabled with rpm=0 stays inactive");
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
temp_env::with_vars(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_ENABLE, Some("false")),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_RPM, Some("600")),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_BURST, None),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
assert!(api_rate_limit_layer_from_env().is_none(), "rpm without enable stays inactive");
|
||||||
|
},
|
||||||
|
);
|
||||||
|
|
||||||
|
temp_env::with_vars(
|
||||||
|
[
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_ENABLE, Some("true")),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_RPM, Some("600")),
|
||||||
|
(rustfs_config::ENV_API_RATE_LIMIT_BURST, Some("50")),
|
||||||
|
],
|
||||||
|
|| {
|
||||||
|
let layer = api_rate_limit_layer_from_env().expect("enabled with rpm");
|
||||||
|
assert_eq!(
|
||||||
|
layer.quota(),
|
||||||
|
RateLimitQuota {
|
||||||
|
requests_per_minute: 600,
|
||||||
|
burst: 50
|
||||||
|
}
|
||||||
|
);
|
||||||
|
},
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -88,7 +88,6 @@ checked_files=(
|
|||||||
"crates/protocols/src/swift/container.rs"
|
"crates/protocols/src/swift/container.rs"
|
||||||
"crates/protocols/src/swift/object.rs"
|
"crates/protocols/src/swift/object.rs"
|
||||||
"crates/protocols/src/swift/formpost.rs"
|
"crates/protocols/src/swift/formpost.rs"
|
||||||
"crates/protocols/src/swift/ratelimit.rs"
|
|
||||||
"crates/protocols/src/swift/expiration.rs"
|
"crates/protocols/src/swift/expiration.rs"
|
||||||
"crates/protocols/src/swift/cors.rs"
|
"crates/protocols/src/swift/cors.rs"
|
||||||
"crates/protocols/src/swift/quota.rs"
|
"crates/protocols/src/swift/quota.rs"
|
||||||
|
|||||||
Reference in New Issue
Block a user