diff --git a/crates/config/src/constants/api.rs b/crates/config/src/constants/api.rs new file mode 100644 index 000000000..86b687c9c --- /dev/null +++ b/crates/config/src/constants/api.rs @@ -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; diff --git a/crates/config/src/constants/mod.rs b/crates/config/src/constants/mod.rs index a7eb8aa1e..ef5b55858 100644 --- a/crates/config/src/constants/mod.rs +++ b/crates/config/src/constants/mod.rs @@ -12,6 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. +pub(crate) mod api; pub(crate) mod app; pub(crate) mod body_limits; pub(crate) mod capacity; diff --git a/crates/config/src/lib.rs b/crates/config/src/lib.rs index a270febe0..da2b59ffe 100644 --- a/crates/config/src/lib.rs +++ b/crates/config/src/lib.rs @@ -15,6 +15,8 @@ #[cfg(feature = "constants")] pub mod constants; #[cfg(feature = "constants")] +pub use constants::api::*; +#[cfg(feature = "constants")] pub use constants::app::*; #[cfg(feature = "constants")] pub use constants::body_limits::*; diff --git a/crates/e2e_test/src/api_rate_limit_test.rs b/crates/e2e_test/src/api_rate_limit_test.rs new file mode 100644 index 000000000..7da385209 --- /dev/null +++ b/crates/e2e_test/src/api_rate_limit_test.rs @@ -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>; + +#[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::().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("TooManyRequests"), "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(()) +} diff --git a/crates/e2e_test/src/lib.rs b/crates/e2e_test/src/lib.rs index eb3e47d70..014e4467f 100644 --- a/crates/e2e_test/src/lib.rs +++ b/crates/e2e_test/src/lib.rs @@ -67,6 +67,10 @@ mod bucket_policy_check_test; #[cfg(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) #[cfg(test)] mod admin_auth_test; diff --git a/crates/protocols/src/swift/mod.rs b/crates/protocols/src/swift/mod.rs index 03f908f79..528cce52b 100644 --- a/crates/protocols/src/swift/mod.rs +++ b/crates/protocols/src/swift/mod.rs @@ -46,7 +46,6 @@ pub mod formpost; pub mod handler; pub mod object; pub mod quota; -pub mod ratelimit; pub mod router; pub mod slo; pub mod staticweb; diff --git a/crates/protocols/src/swift/ratelimit.rs b/crates/protocols/src/swift/ratelimit.rs deleted file mode 100644 index 5c559e4d7..000000000 --- a/crates/protocols/src/swift/ratelimit.rs +++ /dev/null @@ -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 { - 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::() - .map_err(|_| SwiftError::BadRequest(format!("Invalid rate limit value: {}", parts[0])))?; - - let window_seconds = parts[1] - .parse::() - .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 { - // 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>>, -} - -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) -> Option { - // 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"); - } -} diff --git a/crates/protocols/tests/swift_phase4_integration.rs b/crates/protocols/tests/swift_phase4_integration.rs index ab298dbd2..534d13ea1 100644 --- a/crates/protocols/tests/swift_phase4_integration.rs +++ b/crates/protocols/tests/swift_phase4_integration.rs @@ -43,36 +43,4 @@ mod swift_integration { let parsed = expiration::parse_delete_at(delete_at).unwrap(); 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); - } } diff --git a/crates/protocols/tests/swift_simple_integration.rs b/crates/protocols/tests/swift_simple_integration.rs index 783587bc1..02f29c9ef 100644 --- a/crates/protocols/tests/swift_simple_integration.rs +++ b/crates/protocols/tests/swift_simple_integration.rs @@ -16,7 +16,7 @@ #![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; /// Test sync configuration parsing @@ -92,14 +92,6 @@ fn test_symlink_detection() { 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] fn test_quota_structure() { diff --git a/rustfs/src/admin/console.rs b/rustfs/src/admin/console.rs index bf73fa987..5a7a4a1d2 100644 --- a/rustfs/src/admin/console.rs +++ b/rustfs/src/admin/console.rs @@ -16,6 +16,10 @@ use crate::admin::runtime_sources::{current_oidc_handle, default_admin_usecase}; use crate::admin::storage_api::access::RequestContext; use crate::license::has_valid_license; 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::{ CONSOLE_PREFIX, FAVICON_PATH, HEALTH_PREFIX, HEALTH_READY_PATH, HeaderMapCarrier, HealthProbe, LICENSE, RUSTFS_ADMIN_PREFIX, RequestContextLayer, VERSION, build_health_response_parts, collect_probe_readiness, @@ -36,7 +40,7 @@ use rust_embed::RustEmbed; use serde::Serialize; use std::{ net::{IpAddr, SocketAddr}, - sync::OnceLock, + sync::{Arc, OnceLock}, time::Duration, }; use tower_http::catch_panic::CatchPanicLayer; @@ -569,12 +573,38 @@ fn setup_console_middleware_stack( // Add request body limit (10MB for console uploads) .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 { - info!("Console rate limiting enabled: {} requests per minute", rate_limit_rpm); - // Note: tower-http doesn't provide a built-in rate limiter, but we have the foundation - // For production, you would integrate with a rate limiting service like Redis - // For now, we log that it's configured and ready for integration + if let Some(quota) = RateLimitQuota::per_minute(rate_limit_rpm, 0) { + let limiter = Arc::new(RateLimiter::new(quota)); + app = app.layer(middleware::from_fn(move |req: Request, next: middleware::Next| { + 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 @@ -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")] async fn console_trace_layer_records_request_id_on_current_span() { let writer = SharedWriter::default(); diff --git a/rustfs/src/server/http.rs b/rustfs/src/server/http.rs index 9a7ff13d9..7f30ade17 100644 --- a/rustfs/src/server/http.rs +++ b/rustfs/src/server/http.rs @@ -27,6 +27,7 @@ use crate::server::{ RedirectLayer, RequestContextLayer, RequestLoggingLayer, S3ErrorMessageCompatLayer, VirtualHostStyleHintLayer, redact_sensitive_uri_query, }, + rate_limit::{RateLimitLayer, api_rate_limit_layer_from_env}, tls_material::{ 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_HOST_ROUTING: &str = "http_host_routing"; 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_ACCEPT_LOOP_STATE: &str = "http_accept_loop_state"; 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 server_domains_configured = !config.server_domains.is_empty(); let task_handle = tokio::spawn(async move { @@ -890,6 +920,7 @@ pub async fn start_http_server( readiness: readiness.clone(), keystone_auth: auth_keystone::get_keystone_auth(), 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()); @@ -946,6 +977,9 @@ struct ConnectionContext { keystone_auth: Option>, /// Pre-computed trusted proxy layer (avoids per-connection is_enabled() check). trusted_proxy_layer: Option, + /// 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, } #[derive(Clone)] @@ -1130,6 +1164,7 @@ fn process_connection( readiness, keystone_auth, trusted_proxy_layer, + rate_limit_layer, } = context; // Build the hybrid service per-connection. @@ -1183,23 +1218,24 @@ fn process_connection( // 5. RequestContextLayer — creates RequestContext in extensions // 6. EmptyBodyContentLengthCompatLayer — adds Content-Length: 0 for known empty-body API routes // 7. CatchPanicLayer — panic → 500 - // 8. ReadinessGateLayer — blocks until ready - // 9. KeystoneAuthLayer — X-Auth-Token validation - // 10. TraceLayer — request span creation + metrics - // 11. RequestLoggingLayer — single completion event per request - // 12. PropagateRequestIdLayer — X-Request-ID → response - // 13. CompressionLayer — response compression (whitelist, path-aware) - // 14. PathCategoryInjectionLayer — injects path category for compression predicate - // 15. S3ErrorMessageCompatLayer — missing S3 error message compatibility - // 16. IcebergRestErrorCompatLayer — Iceberg REST JSON error compatibility - // 17. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes - // 18. ConditionalCorsLayer — S3 API CORS - // 19. RedirectLayer — console redirect (conditional) - // 20. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses - // 21. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses - // 22. PublicHealthEndpointLayer — handles public health before s3s host parsing - // 23. VirtualHostStyleHintLayer — actionable error for unroutable virtual-hosted-style (conditional) - // 24. DoubleSlashListBucketsCompatLayer — rewrites `GET //` to `GET /` for ListBuckets (MinIO browser compat) + // 8. RateLimitLayer — conditional (external stack only), per-client 429 throttling + // 9. ReadinessGateLayer — blocks until ready + // 10. KeystoneAuthLayer — X-Auth-Token validation + // 11. TraceLayer — request span creation + metrics + // 12. RequestLoggingLayer — single completion event per request + // 13. PropagateRequestIdLayer — X-Request-ID → response + // 14. CompressionLayer — response compression (whitelist, path-aware) + // 15. PathCategoryInjectionLayer — injects path category for compression predicate + // 16. S3ErrorMessageCompatLayer — missing S3 error message compatibility + // 17. IcebergRestErrorCompatLayer — Iceberg REST JSON error compatibility + // 18. ObjectAttributesEtagFixLayer — ETag fix for GetObjectAttributes + // 19. ConditionalCorsLayer — S3 API CORS + // 20. RedirectLayer — console redirect (conditional) + // 21. BodylessStatusFixLayer — clears body for 1xx/204/205/304 responses + // 22. HeadRequestBodyFixLayer — strips actual body bytes from HEAD responses + // 23. PublicHealthEndpointLayer — handles public health before s3s host parsing + // 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 // identical while giving each path family a named construction boundary. @@ -1223,6 +1259,12 @@ fn process_connection( .layer(InternodeRequestContextLiteLayer) .layer(EmptyBodyContentLengthCompatLayer) .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 // This stops requests from hitting IAMAuth or Storage if they are not ready. .layer(ReadinessGateLayer::new(readiness.clone())) diff --git a/rustfs/src/server/rate_limit.rs b/rustfs/src/server/rate_limit.rs index db1c0cf47..4790e24af 100644 --- a/rustfs/src/server/rate_limit.rs +++ b/rustfs/src/server/rate_limit.rs @@ -12,124 +12,391 @@ // See the License for the specific language governing permissions and // 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. -//! It helps protect against DoS attacks by limiting the number of requests per client. +//! A token-bucket limiter keyed by client IP. Over-limit requests are rejected +//! 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::hash::{BuildHasher, RandomState}; use std::net::IpAddr; -use std::sync::{Arc, RwLock}; -use std::time::{Duration, Instant}; +use std::sync::{Arc, Mutex, PoisonError}; +use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; use tower::{Layer, Service}; +use tracing::{debug, warn}; -/// Configuration for rate limiting -#[derive(Debug, Clone)] -pub struct RateLimitConfig { - /// Maximum number of requests per window - pub max_requests: u32, - /// Time window duration - pub window_duration: Duration, +const LOG_COMPONENT_SERVER: &str = "server"; +const LOG_SUBSYSTEM_RATE_LIMIT: &str = "rate_limit"; +const EVENT_REQUEST_RATE_LIMITED: &str = "request_rate_limited"; + +pub(crate) const METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL: &str = "rustfs_http_server_requests_rate_limited_total"; +pub(crate) const LABEL_RATE_LIMIT_SCOPE: &str = "scope"; +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 { - fn default() -> Self { - Self { - max_requests: 100, - window_duration: Duration::from_secs(60), +impl RateLimitQuota { + /// Build a quota from operator configuration. `rpm == 0` means unlimited + /// (`None`); `burst == 0` means "same as RPM". + pub fn per_minute(rpm: u32, burst: u32) -> Option { + 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)] struct TokenBucket { - /// Maximum number of tokens - max_tokens: u32, - /// Current number of tokens tokens: f64, - /// Rate at which tokens are refilled (tokens per second) - refill_rate: f64, - /// Last time tokens were refilled last_refill: Instant, } impl TokenBucket { - fn new(max_tokens: u32, window_duration: Duration) -> Self { + fn full(capacity: f64, now: Instant) -> Self { Self { - max_tokens, - tokens: max_tokens as f64, - refill_rate: max_tokens as f64 / window_duration.as_secs_f64(), - last_refill: Instant::now(), + tokens: capacity, + last_refill: now, } } - fn try_consume(&mut self) -> bool { - self.refill(); + fn try_consume(&mut self, capacity: f64, refill_per_sec: f64, now: Instant) -> RateLimitDecision { + 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 { self.tokens -= 1.0; - true + RateLimitDecision::Allowed } 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) { - let now = Instant::now(); - let elapsed = now.duration_since(self.last_refill).as_secs_f64(); - self.tokens = (self.tokens + elapsed * self.refill_rate).min(self.max_tokens as f64); - self.last_refill = now; + /// Idle long enough that refill would have restored a full bucket: the + /// entry carries no information a fresh one would not, so it can be + /// dropped losslessly. + fn is_refilled_idle(&self, capacity: f64, refill_per_sec: f64, now: Instant) -> bool { + now.saturating_duration_since(self.last_refill).as_secs_f64() * refill_per_sec >= capacity } } -/// Rate limiter state #[derive(Debug)] -struct RateLimiterState { - /// Per-client token buckets +struct Shard { clients: HashMap, - /// Configuration - config: RateLimitConfig, + next_cleanup: Instant, } -impl RateLimiterState { - fn new(config: RateLimitConfig) -> Self { +/// Sharded per-IP token-bucket limiter. Shared by the S3 API tower layer and +/// the console middleware (each with its own instance and quota). +#[derive(Debug)] +pub struct RateLimiter { + shards: Vec>, + 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 { - clients: HashMap::new(), - config, + shards: (0..shard_count.max(1)) + .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 { - let bucket = self - .clients - .entry(client_ip) - .or_insert_with(|| TokenBucket::new(self.config.max_requests, self.config.window_duration)); - - bucket.try_consume() + pub fn quota(&self) -> RateLimitQuota { + self.quota } - fn cleanup_expired(&mut self) { - let now = Instant::now(); - let window = self.config.window_duration; - self.clients - .retain(|_, bucket| now.duration_since(bucket.last_refill) < window * 2); + /// Consume one token for `ip`, creating its bucket on first sight. + pub fn check(&self, ip: IpAddr) -> RateLimitDecision { + self.check_at(ip, Instant::now()) + } + + 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(req: &Request) -> Option { + if let Some(info) = req.extensions().get::() { + return Some(info.real_ip); + } + req.extensions().get::().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; +type BoxBody = http_body_util::combinators::UnsyncBoxBody; + +/// 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 { + // 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!("{id}")) + .unwrap_or_default(); + let body = format!( + "\ + TooManyRequests\ + Request rate limit exceeded. Reduce your request rate and retry after the indicated delay.\ + {request_id_xml}" + ); + 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 { + 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)] pub struct RateLimitLayer { - state: Arc>, + limiter: Arc, } impl RateLimitLayer { - pub fn new(config: RateLimitConfig) -> Self { + pub fn new(quota: RateLimitQuota) -> 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 Layer for RateLimitLayer { @@ -138,75 +405,401 @@ impl Layer for RateLimitLayer { fn layer(&self, inner: S) -> Self::Service { RateLimitService { 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)] pub struct RateLimitService { inner: S, - state: Arc>, + limiter: Arc, } -impl Service> for RateLimitService +impl Service> for RateLimitService where - S: Service, Response = Response>, + S: Service, Response = Response>, { - type Response = S::Response; + type Response = Response; type Error = S::Error; - type Future = S::Future; + type Future = Either, S::Error>>>; fn poll_ready(&mut self, cx: &mut std::task::Context<'_>) -> std::task::Poll> { self.inner.poll_ready(cx) } fn call(&mut self, req: Request) -> Self::Future { - // Extract client IP from request - let client_ip = extract_client_ip(&req); - - // Check rate limit - let allowed = { - let mut state = self.state.write().unwrap_or_else(|e| e.into_inner()); - // Periodically cleanup expired entries - if rand::random::().is_multiple_of(100) { - state.cleanup_expired(); - } - state.check_rate_limit(client_ip) + if is_rate_limit_exempt_path(req.uri().path()) { + return Either::Left(self.inner.call(req)); + } + let Some(ip) = client_ip(&req) else { + // Fail open: without a socket peer address there is no trustworthy + // client identity, and deriving one from spoofable headers would + // let clients pick their own bucket. + return Either::Left(self.inner.call(req)); }; - if allowed { - self.inner.call(req) - } else { - // Return 429 Too Many Requests - // Note: In a real implementation, you'd want to return a proper 429 response. - // For now, we just pass the request through to the inner service. - // The rate limiting is enforced at the application level. - self.inner.call(req) + match self.limiter.check(ip) { + RateLimitDecision::Allowed => Either::Left(self.inner.call(req)), + RateLimitDecision::Limited(throttle) => { + counter!( + METRIC_HTTP_SERVER_REQUESTS_RATE_LIMITED_TOTAL, + LABEL_RATE_LIMIT_SCOPE => RATE_LIMIT_SCOPE_S3_API + ) + .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 -fn extract_client_ip(req: &Request) -> IpAddr { - // Try X-Forwarded-For header first - if let Some(forwarded) = req.headers().get("x-forwarded-for") - && let Ok(forwarded_str) = forwarded.to_str() - && let Some(first_ip) = forwarded_str.split(',').next() - && let Ok(ip) = first_ip.trim().parse::() - { - return ip; +#[cfg(test)] +mod tests { + use super::*; + use http_body_util::BodyExt; + use serial_test::serial; + use std::net::{Ipv4Addr, SocketAddr}; + use std::task::{Context, Poll}; + + fn quota(rpm: u32, burst: u32) -> RateLimitQuota { + RateLimitQuota::per_minute(rpm, burst).expect("non-zero rpm") } - // Try X-Real-IP header - if let Some(real_ip) = req.headers().get("x-real-ip") - && let Ok(ip_str) = real_ip.to_str() - && let Ok(ip) = ip_str.parse::() - { - return ip; + fn ip(last: u8) -> IpAddr { + IpAddr::V4(Ipv4Addr::new(203, 0, 113, last)) } - // Default to localhost - IpAddr::V4(std::net::Ipv4Addr::LOCALHOST) + #[test] + 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 Service> for OkService { + type Response = Response; + type Error = std::convert::Infallible; + type Future = Ready>; + + fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn call(&mut self, _req: Request) -> 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 { + 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("TooManyRequests"), "body: {body}"); + assert!(body.contains("req-123"), "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("evil")); + 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(""), "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 + } + ); + }, + ); + } } diff --git a/scripts/check_logging_guardrails.sh b/scripts/check_logging_guardrails.sh index 088e788ee..cb3640e5e 100755 --- a/scripts/check_logging_guardrails.sh +++ b/scripts/check_logging_guardrails.sh @@ -88,7 +88,6 @@ checked_files=( "crates/protocols/src/swift/container.rs" "crates/protocols/src/swift/object.rs" "crates/protocols/src/swift/formpost.rs" - "crates/protocols/src/swift/ratelimit.rs" "crates/protocols/src/swift/expiration.rs" "crates/protocols/src/swift/cors.rs" "crates/protocols/src/swift/quota.rs"