diff --git a/crates/buzz-auth/src/rate_limit.rs b/crates/buzz-auth/src/rate_limit.rs index 8fd42c50fb..7b7ecd0236 100644 --- a/crates/buzz-auth/src/rate_limit.rs +++ b/crates/buzz-auth/src/rate_limit.rs @@ -76,6 +76,16 @@ impl LimitType { Self::IpConnections => "conn", } } + + /// Human-readable name used in rejection messages and metric labels. + pub fn as_str(&self) -> &'static str { + match self { + Self::Messages => "messages", + Self::ApiCalls => "api_calls", + Self::WsEvents => "ws_events", + Self::IpConnections => "ip_connections", + } + } } /// Per-tier rate limit thresholds. @@ -96,15 +106,13 @@ pub struct RateLimitConfig { /// Maximum messages per minute for standard-tier agent tokens. Default: 120. #[serde(default = "default_agent_std_msg")] pub agent_standard_messages_per_min: u64, - /// Maximum HTTP API calls per minute for standard-tier agent tokens. Default: 600. - #[serde(default = "default_agent_std_api")] - pub agent_standard_api_calls_per_min: u64, - /// Maximum messages per minute for elevated-tier agent tokens. Default: 300. - #[serde(default = "default_agent_elev_msg")] - pub agent_elevated_messages_per_min: u64, - /// Maximum messages per minute for platform-tier agent tokens. Default: 600. - #[serde(default = "default_agent_plat_msg")] - pub agent_platform_messages_per_min: u64, + /// Maximum WebSocket events per second for agent tokens. Default: 10. + /// + /// Agents inherit the same burst ceiling as humans by default. Tune via + /// `BUZZ_RATE_LIMIT_AGENT_WS_EVENTS_PER_SEC` once `limit_type` instrumentation + /// data establishes the right operating value. + #[serde(default = "default_agent_ws")] + pub agent_ws_events_per_sec: u64, } fn default_human_msg() -> u64 { @@ -119,14 +127,10 @@ fn default_human_ws() -> u64 { fn default_agent_std_msg() -> u64 { 120 } -fn default_agent_std_api() -> u64 { - 600 -} -fn default_agent_elev_msg() -> u64 { - 300 -} -fn default_agent_plat_msg() -> u64 { - 600 +fn default_agent_ws() -> u64 { + // Same as human default: behavior-neutral at merge; tune on builderlab + // once limit_type data from instrumented deployments is available. + 10 } impl Default for RateLimitConfig { @@ -136,9 +140,7 @@ impl Default for RateLimitConfig { human_api_calls_per_min: default_human_api(), human_ws_events_per_sec: default_human_ws(), agent_standard_messages_per_min: default_agent_std_msg(), - agent_standard_api_calls_per_min: default_agent_std_api(), - agent_elevated_messages_per_min: default_agent_elev_msg(), - agent_platform_messages_per_min: default_agent_plat_msg(), + agent_ws_events_per_sec: default_agent_ws(), } } } @@ -323,4 +325,21 @@ mod tests { .unwrap(); assert!(result.allowed); } + + #[test] + fn limit_type_as_str_is_stable() { + // These strings appear in relay NOTICE text and metric labels; + // changing them is a breaking observability change. + assert_eq!(LimitType::Messages.as_str(), "messages"); + assert_eq!(LimitType::ApiCalls.as_str(), "api_calls"); + assert_eq!(LimitType::WsEvents.as_str(), "ws_events"); + assert_eq!(LimitType::IpConnections.as_str(), "ip_connections"); + } + + #[test] + fn rate_limit_config_default_has_agent_ws_field() { + let cfg = RateLimitConfig::default(); + // agent_ws_events_per_sec starts equal to human to be behavior-neutral at merge. + assert_eq!(cfg.agent_ws_events_per_sec, cfg.human_ws_events_per_sec); + } } diff --git a/crates/buzz-relay/src/api/bridge.rs b/crates/buzz-relay/src/api/bridge.rs index a118ff453f..cd3985f1e1 100644 --- a/crates/buzz-relay/src/api/bridge.rs +++ b/crates/buzz-relay/src/api/bridge.rs @@ -39,14 +39,29 @@ async fn enforce_http_admission( { Ok(()) => Ok(()), Err(crate::admission::AdmissionError::Exceeded { reset_in_secs }) => { - metrics::counter!("buzz_admission_rejections_total", "transport" => "http", "reason" => "quota").increment(1); + metrics::counter!( + "buzz_admission_rejections_total", + "transport" => "http", + "reason" => "quota", + "limit_type" => LimitType::ApiCalls.as_str(), + ) + .increment(1); Err(api_error( StatusCode::TOO_MANY_REQUESTS, - &format!("rate-limited: quota exceeded; retry in {reset_in_secs}s"), + &format!( + "rate-limited: quota exceeded ({}); retry in {reset_in_secs}s", + LimitType::ApiCalls.as_str() + ), )) } Err(crate::admission::AdmissionError::Unavailable) => { - metrics::counter!("buzz_admission_rejections_total", "transport" => "http", "reason" => "unavailable").increment(1); + metrics::counter!( + "buzz_admission_rejections_total", + "transport" => "http", + "reason" => "unavailable", + "limit_type" => LimitType::ApiCalls.as_str(), + ) + .increment(1); Err(api_error( StatusCode::SERVICE_UNAVAILABLE, "rate-limited: shared admission unavailable", diff --git a/crates/buzz-relay/src/config.rs b/crates/buzz-relay/src/config.rs index 85a0ca2efe..a930aa6bb2 100644 --- a/crates/buzz-relay/src/config.rs +++ b/crates/buzz-relay/src/config.rs @@ -317,17 +317,9 @@ fn rate_limit_config_from_env() -> Result Some(sub_id.as_str()), _ => None, }; - if !send_admission_result(conn, ws_result, sub_id) { + if !send_admission_result(conn, ws_result, sub_id, LimitType::WsEvents) { return false; } if is_event { - let message_limit = if is_agent { - limits.agent_standard_messages_per_min - } else { - limits.human_messages_per_min - }; - let message_result = crate::admission::check_principal( - state.admission_rate_limiter.as_ref(), - &conn.tenant, - &pubkey, - LimitType::Messages, - 60, - message_limit, - ) - .await; - if !send_admission_result(conn, message_result, None) { - return false; + // Ephemeral events (kinds 20000–29999) are never persisted; billing them + // against the durable-message quota allows telemetry (observer frames, + // typing indicators, presence) to starve real traffic. They still count + // against WsEvents above, so the relay's per-second burst protection holds. + let is_ephemeral_event = matches!(msg, ClientMessage::Event(e) if is_ephemeral(buzz_core::kind::event_kind_u32(e))); + if !is_ephemeral_event { + let message_limit = if is_agent { + limits.agent_standard_messages_per_min + } else { + limits.human_messages_per_min + }; + let message_result = crate::admission::check_principal( + state.admission_rate_limiter.as_ref(), + &conn.tenant, + &pubkey, + LimitType::Messages, + 60, + message_limit, + ) + .await; + if !send_admission_result(conn, message_result, None, LimitType::Messages) { + return false; + } } } @@ -656,19 +667,35 @@ fn send_admission_result( conn: &ConnectionState, result: Result<(), crate::admission::AdmissionError>, sub_id: Option<&str>, + limit_type: LimitType, ) -> bool { match result { Ok(()) => true, Err(crate::admission::AdmissionError::Exceeded { reset_in_secs }) => { - metrics::counter!("buzz_admission_rejections_total", "transport" => "websocket", "reason" => "quota").increment(1); + metrics::counter!( + "buzz_admission_rejections_total", + "transport" => "websocket", + "reason" => "quota", + "limit_type" => limit_type.as_str(), + ) + .increment(1); conn.send(request_rejection_message( sub_id, - &format!("rate-limited: quota exceeded; retry in {reset_in_secs}s"), + &format!( + "rate-limited: quota exceeded ({}); retry in {reset_in_secs}s", + limit_type.as_str() + ), )); false } Err(crate::admission::AdmissionError::Unavailable) => { - metrics::counter!("buzz_admission_rejections_total", "transport" => "websocket", "reason" => "unavailable").increment(1); + metrics::counter!( + "buzz_admission_rejections_total", + "transport" => "websocket", + "reason" => "unavailable", + "limit_type" => limit_type.as_str(), + ) + .increment(1); conn.send(request_rejection_message( sub_id, "rate-limited: shared admission unavailable", @@ -785,6 +812,108 @@ mod tests { assert_eq!(notice, serde_json::json!(["NOTICE", reason])); } + /// Build a minimal ConnectionState whose `send_tx` is readable in tests. + fn test_connection() -> (ConnectionState, mpsc::Receiver) { + let (send_tx, recv) = mpsc::channel(16); + let (ctrl_tx, _ctrl_rx) = mpsc::channel(1); + let conn = ConnectionState { + conn_id: Uuid::nil(), + tenant: buzz_core::tenant::TenantContext::resolved( + buzz_core::tenant::CommunityId::from_uuid(Uuid::nil()), + "test.local".to_string(), + ), + remote_addr: "127.0.0.1:9999".parse().unwrap(), + auth_state: RwLock::new(AuthState::Failed), + subscriptions: Arc::new(tokio::sync::Mutex::new(HashMap::new())), + send_tx, + ctrl_tx, + cancel: CancellationToken::new(), + backpressure_count: Arc::new(AtomicU8::new(0)), + grace_limit: 10, + }; + (conn, recv) + } + + #[test] + fn send_admission_result_exceeded_notice_names_limit_type_and_keeps_retry_phrase() { + let (conn, mut rx) = test_connection(); + let err = crate::admission::AdmissionError::Exceeded { reset_in_secs: 42 }; + + let result = send_admission_result(&conn, Err(err), None, LimitType::Messages); + + assert!(!result, "exceeded should return false"); + let msg = rx.try_recv().expect("should have sent one message"); + let text = match msg { + WsMessage::Text(t) => t.to_string(), + other => panic!("unexpected frame: {other:?}"), + }; + let parsed: serde_json::Value = serde_json::from_str(&text).expect("valid json"); + let notice_text = parsed[1].as_str().expect("NOTICE text"); + assert!( + notice_text.contains("messages"), + "NOTICE must name the limit type: {notice_text}" + ); + assert!( + notice_text.contains("retry in 42s"), + "NOTICE must preserve 'retry in Ns' phrase for client parsers: {notice_text}" + ); + } + + #[test] + fn send_admission_result_exceeded_ws_events_names_ws_events() { + let (conn, mut rx) = test_connection(); + let err = crate::admission::AdmissionError::Exceeded { reset_in_secs: 3 }; + + send_admission_result(&conn, Err(err), None, LimitType::WsEvents); + + let msg = rx.try_recv().expect("should have sent one message"); + let text = match msg { + WsMessage::Text(t) => t.to_string(), + other => panic!("unexpected frame: {other:?}"), + }; + let parsed: serde_json::Value = serde_json::from_str(&text).expect("valid json"); + let notice_text = parsed[1].as_str().expect("NOTICE text"); + assert!( + notice_text.contains("ws_events"), + "NOTICE for WsEvents must say 'ws_events': {notice_text}" + ); + assert!( + notice_text.contains("retry in 3s"), + "NOTICE must preserve 'retry in Ns' phrase: {notice_text}" + ); + } + + #[test] + fn send_admission_result_exceeded_with_sub_id_emits_closed_not_notice() { + let (conn, mut rx) = test_connection(); + let err = crate::admission::AdmissionError::Exceeded { reset_in_secs: 5 }; + + send_admission_result(&conn, Err(err), Some("sub-xyz"), LimitType::WsEvents); + + let msg = rx.try_recv().expect("should have sent one message"); + let text = match msg { + WsMessage::Text(t) => t.to_string(), + other => panic!("unexpected frame: {other:?}"), + }; + let parsed: serde_json::Value = serde_json::from_str(&text).expect("valid json"); + assert_eq!( + parsed[0].as_str(), + Some("CLOSED"), + "sub-scoped rejection should emit CLOSED" + ); + assert_eq!(parsed[1].as_str(), Some("sub-xyz")); + } + + #[test] + fn send_admission_result_ok_returns_true_and_sends_nothing() { + let (conn, mut rx) = test_connection(); + + let result = send_admission_result(&conn, Ok(()), None, LimitType::Messages); + + assert!(result, "allowed result should return true"); + assert!(rx.try_recv().is_err(), "no message should be sent on Ok"); + } + #[tokio::test] async fn send_loop_batches_queued_data_frames_into_one_flush() { let (data_tx, data_rx) = mpsc::channel(MAX_WS_SEND_BATCH);