diff --git a/helm/values.yaml b/helm/values.yaml index 2524ad9..865062a 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -28,7 +28,7 @@ serviceAccount: # If not set and create is true, a name is generated using the fullname template name: "" # This is a list of cluster permissions to apply to the service account. - # By default it grants all permissions. + # The default is what robotlb needs; a custom list replaces it entirely. permissions: - apiGroups: [""] resources: [services, services/status] @@ -39,6 +39,9 @@ serviceAccount: - apiGroups: [discovery.k8s.io] resources: [endpointslices] verbs: [get, list, watch] + - apiGroups: [events.k8s.io] + resources: [events] + verbs: [create] podAnnotations: {} podLabels: {} diff --git a/src/error.rs b/src/error.rs index dae2e3f..ec91801 100644 --- a/src/error.rs +++ b/src/error.rs @@ -22,65 +22,229 @@ pub enum RobotLBError { UnknownLBAlgorithm, #[error("Cannot get target nodes, because the service has no selector")] ServiceWithoutSelector, + #[error("Hetzner Cloud API rate limit reached, the pause ends in {}s", .0.as_millis().div_ceil(1000))] + RateLimited(std::time::Duration), // HCloud API errors - #[error("Cannot attach load balancer to a network. Reason: {0}")] + #[error("Cannot attach load balancer to a network. Reason: {}", describe(.0))] HCloudLBAttachToNetworkError( #[from] hcloud::apis::Error, ), - #[error("Cannot detach load balancer from network. Reason: {0}")] + #[error("Cannot detach load balancer from network. Reason: {}", describe(.0))] HcloudLBDetachFromNetworkError( #[from] hcloud::apis::Error, ), - #[error("Cannot add load balancer target. Reason: {0}")] + #[error("Cannot add load balancer target. Reason: {}", describe(.0))] HcloudLBAddTargetError( #[from] hcloud::apis::Error, ), - #[error("Cannot remove load balancer target. Reason: {0}")] + #[error("Cannot remove load balancer target. Reason: {}", describe(.0))] HcloudLBRemoveTargetError( #[from] hcloud::apis::Error, ), - #[error("Cannot add service to load balancer. Reason: {0}")] + #[error("Cannot add service to load balancer. Reason: {}", describe(.0))] HcloudLBAddServiceError( #[from] hcloud::apis::Error, ), - #[error("Cannot remove service from load balancer. Reason: {0}")] + #[error("Cannot remove service from load balancer. Reason: {}", describe(.0))] HcloudLBRemoveServiceError( #[from] hcloud::apis::Error, ), - #[error("Cannot create load balancer. Reason: {0}")] + #[error("Cannot create load balancer. Reason: {}", describe(.0))] HcloudLBCreateError( #[from] hcloud::apis::Error, ), - #[error("Cannot delete load balancer. Reason: {0}")] + #[error("Cannot delete load balancer. Reason: {}", describe(.0))] HcloudLBDeleteError( #[from] hcloud::apis::Error, ), - #[error("Cannot get load balancer. Reason: {0}")] + #[error("Cannot get load balancer. Reason: {}", describe(.0))] HcloudLBGetError( #[from] hcloud::apis::Error, ), - #[error("Cannot update service. Reason: {0}")] + #[error("Cannot update service. Reason: {}", describe(.0))] HcloudLBUpdateServiceError( #[from] hcloud::apis::Error, ), - #[error("Cannot change type of load balancer. Reason: {0}")] + #[error("Cannot change type of load balancer. Reason: {}", describe(.0))] HcloudLBChangeType( #[from] hcloud::apis::Error, ), - #[error("Cannot change algorithm of load balancer. Reason: {0}")] + #[error("Cannot change algorithm of load balancer. Reason: {}", describe(.0))] HcloudLBChangeAlgorithm( #[from] hcloud::apis::Error, ), - #[error("Cannot list networks. Reason: {0}")] + #[error("Cannot list networks. Reason: {}", describe(.0))] HcloudListNetworksError( #[from] hcloud::apis::Error, ), - #[error("Cannot list load balancers. Reason: {0}")] + #[error("Cannot list load balancers. Reason: {}", describe(.0))] HcloudListLoadBalancersError( #[from] hcloud::apis::Error, ), } + +impl RobotLBError { + /// Whether Hetzner rejected the call because the project ran out of API requests. + #[must_use] + pub fn is_rate_limited(&self) -> bool { + // No wildcard arm: a new variant must be sorted into one of the two groups. + match self { + Self::HCloudLBAttachToNetworkError(error) => is_rate_limit_response(error), + Self::HcloudLBDetachFromNetworkError(error) => is_rate_limit_response(error), + Self::HcloudLBAddTargetError(error) => is_rate_limit_response(error), + Self::HcloudLBRemoveTargetError(error) => is_rate_limit_response(error), + Self::HcloudLBAddServiceError(error) => is_rate_limit_response(error), + Self::HcloudLBRemoveServiceError(error) => is_rate_limit_response(error), + Self::HcloudLBCreateError(error) => is_rate_limit_response(error), + Self::HcloudLBDeleteError(error) => is_rate_limit_response(error), + Self::HcloudLBGetError(error) => is_rate_limit_response(error), + Self::HcloudLBUpdateServiceError(error) => is_rate_limit_response(error), + Self::HcloudLBChangeType(error) => is_rate_limit_response(error), + Self::HcloudLBChangeAlgorithm(error) => is_rate_limit_response(error), + Self::HcloudListNetworksError(error) => is_rate_limit_response(error), + Self::HcloudListLoadBalancersError(error) => is_rate_limit_response(error), + Self::InvalidNodeFilter(_) + | Self::UnsupportedServiceType + | Self::SkipService + | Self::PaseIntError(_) + | Self::PaseBoolError(_) + | Self::HCloudError(_) + | Self::KubeError(_) + | Self::UnknownLBAlgorithm + | Self::ServiceWithoutSelector + | Self::RateLimited(_) => false, + } + } +} + +/// Whether Hetzner answered 429 because the project ran out of API requests. +#[must_use] +pub fn is_rate_limit_response(error: &hcloud::apis::Error) -> bool { + matches!(error, hcloud::apis::Error::ResponseError(response) if response.status.as_u16() == 429) +} + +/// One line with the HTTP status and the error Hetzner reported, if the body carries one. +#[must_use] +pub fn describe(error: &hcloud::apis::Error) -> String { + let hcloud::apis::Error::ResponseError(response) = error else { + return error.to_string(); + }; + let body = + k8s_openapi::serde_json::from_str::(&response.content).ok(); + let reported = body.as_ref().and_then(|body| { + let error = body.get("error")?; + Some(format!( + "{}: {}", + error.get("code")?.as_str()?, + error.get("message")?.as_str()? + )) + }); + reported.map_or_else( + || response.status.to_string(), + |reported| { + let one_line = reported.split_whitespace().collect::>().join(" "); + format!("{}: {one_line}", response.status) + }, + ) +} + +/// Replace the API token, or any prefix of it long enough to identify it, with a marker. +/// Hetzner quotes the start of the token in some error messages. +#[must_use] +pub fn redact(message: &str, token: &str) -> String { + const SHORTEST_PREFIX: usize = 8; + let mut redacted = message.to_string(); + for len in (SHORTEST_PREFIX..=token.len()).rev() { + if let Some(prefix) = token.get(..len) { + redacted = redacted.replace(prefix, "[REDACTED]"); + } + } + redacted +} + +#[cfg(test)] +mod tests { + use super::{describe, is_rate_limit_response, redact, RobotLBError}; + use hcloud::apis::{load_balancers_api::ListLoadBalancersError, Error, ResponseContent}; + + const TOKEN: &str = "abcdefghijklmnopqrstuvwxyz0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZ01"; + + fn response_error(status: u16, content: &str) -> Error { + Error::ResponseError(ResponseContent { + status: status.try_into().unwrap(), + content: content.to_string(), + entity: None, + }) + } + + #[test] + fn a_token_prefix_is_redacted() { + let message = format!("limit reached for token {}", &TOKEN[..32]); + assert_eq!( + redact(&message, TOKEN), + "limit reached for token [REDACTED]" + ); + } + + #[test] + fn the_whole_token_is_redacted() { + assert_eq!(redact(&format!("x{TOKEN}y"), TOKEN), "x[REDACTED]y"); + } + + #[test] + fn a_message_without_the_token_is_kept() { + assert_eq!(redact("status 500", TOKEN), "status 500"); + } + + #[test] + fn an_empty_token_redacts_nothing() { + assert_eq!(redact("status 500", ""), "status 500"); + } + + #[test] + fn the_hetzner_message_is_described_on_one_line() { + let error = response_error( + 429, + r#"{"error": {"code": "rate_limit_exceeded", "message": "limit\n reached"}}"#, + ); + assert_eq!( + describe(&error), + "429 Too Many Requests: rate_limit_exceeded: limit reached" + ); + } + + #[test] + fn a_body_that_is_not_json_falls_back_to_the_status() { + let error = response_error(502, "bad gateway"); + assert_eq!(describe(&error), "502 Bad Gateway"); + } + + #[test] + fn a_pause_shorter_than_a_second_is_not_reported_as_over() { + let error = RobotLBError::RateLimited(std::time::Duration::from_millis(300)); + assert!(error.to_string().ends_with("in 1s")); + } + + #[test] + fn a_429_is_a_rate_limit() { + let error = RobotLBError::from(response_error(429, "")); + assert!(error.is_rate_limited()); + } + + #[test] + fn only_a_429_response_is_a_rate_limit() { + assert!(is_rate_limit_response(&response_error(429, ""))); + // Hetzner answers 422 for a target outside the vSwitch subnet. + assert!(!is_rate_limit_response(&response_error(422, ""))); + } + + #[test] + fn other_statuses_are_not_a_rate_limit() { + assert!(!RobotLBError::from(response_error(500, "")).is_rate_limited()); + assert!(!RobotLBError::SkipService.is_rate_limited()); + } +} diff --git a/src/lb.rs b/src/lb.rs index 8725314..e1cc5f3 100644 --- a/src/lb.rs +++ b/src/lb.rs @@ -366,9 +366,13 @@ impl LoadBalancer { // which must not keep the remaining nodes out of the load balancer. match added { Ok(_) => live += 1, + // Every further call would be rejected too and only drain the budget. + Err(error) if crate::error::is_rate_limit_response(&error) => { + return Err(error.into()); + } Err(error) => { tracing::warn!("Cannot add target {ip}: {error}"); - last_error = Some(error.to_string()); + last_error = Some(crate::error::describe(&error)); } } } @@ -619,16 +623,9 @@ impl LoadBalancer { }), }, ) - .await; - if let Err(e) = response { - tracing::error!("Failed to create load balancer: {:?}", e); - return Err(RobotLBError::HCloudError(format!( - "Failed to create load balancer: {:?}", - e - ))); - } + .await?; - Ok(*response.unwrap().load_balancer) + Ok(*response.load_balancer) } /// Get the network from Hetzner Cloud. diff --git a/src/main.rs b/src/main.rs index 164cc6f..edcfa73 100644 --- a/src/main.rs +++ b/src/main.rs @@ -18,7 +18,7 @@ use clap::Parser; use config::OperatorConfig; -use error::{RobotLBError, RobotLBResult}; +use error::{redact, RobotLBError, RobotLBResult}; use futures::StreamExt; use hcloud::apis::configuration::Configuration as HCloudConfig; use k8s_openapi::{ @@ -30,12 +30,22 @@ use k8s_openapi::{ }; use kube::{ api::{ListParams, PatchParams}, - runtime::{controller::Action, watcher, Controller}, + runtime::{ + controller::{self, Action}, + events::{Event, EventType, Recorder, Reporter}, + watcher, Controller, + }, Resource, ResourceExt, }; use label_filter::LabelFilter; use lb::{LBService, LoadBalancer}; -use std::{collections::HashSet, str::FromStr, sync::Arc, time::Duration}; +use rate_limit::{spread, RateLimitGate}; +use std::{ + collections::HashSet, + str::FromStr, + sync::Arc, + time::{Duration, Instant}, +}; pub mod config; pub mod consts; @@ -43,6 +53,7 @@ pub mod error; pub mod finalizers; pub mod label_filter; pub mod lb; +pub mod rate_limit; #[cfg(not(target_env = "msvc"))] #[global_allocator] @@ -69,27 +80,38 @@ async fn main() -> RobotLBResult<()> { hcloud_conf, )); tracing::info!("Starting the controller"); + let token = operator_config.hcloud_token; Controller::new( kube::Api::::all(kube_client), watcher::Config::default(), ) .run(reconcile_service, on_error, context) - .for_each(|reconcilation_result| async move { - match reconcilation_result { - Ok((service, _action)) => { - tracing::info!("Reconcilation of a service {} was successful", service.name); - } - Err(err) => match err { + .for_each(|reconcilation_result| { + let token = token.clone(); + async move { + match reconcilation_result { + Ok((service, _action)) => { + tracing::info!("Reconcilation of a service {} was successful", service.name); + } // During reconcilation process, // the controller has decided to skip the service. - kube::runtime::controller::Error::ReconcilerFailed( - RobotLBError::SkipService, - _, - ) => {} - _ => { - tracing::error!("Error reconciling service: {:#?}", err); + Err(controller::Error::ReconcilerFailed(RobotLBError::SkipService, _)) => {} + Err(controller::Error::ReconcilerFailed( + error @ RobotLBError::RateLimited(_), + service, + )) => { + tracing::info!("Service {service}: {error}"); } - }, + Err(controller::Error::ReconcilerFailed(error, service)) => { + tracing::error!( + "Reconcilation of service {service} failed: {}", + redact(&error.to_string(), &token) + ); + } + Err(error) => { + tracing::error!("Controller error: {}", redact(&error.to_string(), &token)); + } + } } }) .await; @@ -101,18 +123,16 @@ pub struct CurrentContext { pub client: kube::Client, pub config: OperatorConfig, pub hcloud_config: HCloudConfig, + pub rate_limit: Arc, } impl CurrentContext { #[must_use] - pub const fn new( - client: kube::Client, - config: OperatorConfig, - hcloud_config: HCloudConfig, - ) -> Self { + pub fn new(client: kube::Client, config: OperatorConfig, hcloud_config: HCloudConfig) -> Self { Self { client, config, hcloud_config, + rate_limit: Arc::default(), } } } @@ -126,6 +146,60 @@ pub async fn reconcile_service( svc: Arc, context: Arc, ) -> RobotLBResult { + let result = sync_service(svc.clone(), context.clone()).await; + if let Err(error) = &result { + if publishes_event(error) { + report_failure(&svc, &context, error).await; + } + } + result +} + +/// Skipped services are every service robotlb does not own. A service waiting at the +/// rate limit gate still gets an event each time it wakes up to a closed gate. +const fn publishes_event(error: &RobotLBError) -> bool { + !matches!(error, RobotLBError::SkipService) +} + +/// Put the error on the service as a warning event, where `kubectl describe` shows it. +async fn report_failure(svc: &Service, context: &CurrentContext, error: &RobotLBError) { + let recorder = Recorder::new( + context.client.clone(), + Reporter { + controller: "robotlb".to_string(), + instance: None, + }, + svc.object_ref(&()), + ); + let published = recorder + .publish(Event { + type_: EventType::Warning, + reason: "SyncLoadBalancerFailed".to_string(), + note: Some(event_note(error, &context.config.hcloud_token)), + action: "Reconcile".to_string(), + secondary: None, + }) + .await; + if let Err(publish_error) = published { + tracing::warn!("Cannot publish an event for the service: {publish_error}"); + } +} + +/// The error as an event note: token redacted and cut to the size the API accepts. +fn event_note(error: &RobotLBError, token: &str) -> String { + const MAX_NOTE_BYTES: usize = 1024; + let mut note = redact(&error.to_string(), token); + if note.len() > MAX_NOTE_BYTES { + let mut end = MAX_NOTE_BYTES; + while !note.is_char_boundary(end) { + end -= 1; + } + note.truncate(end); + } + note +} + +async fn sync_service(svc: Arc, context: Arc) -> RobotLBResult { let svc_type = svc .spec .as_ref() @@ -148,6 +222,12 @@ pub async fn reconcile_service( return Err(RobotLBError::SkipService); } + // Hetzner counts requests per project, so while one service is rate limited + // every other one waits too instead of spending the budget being waited for. + if let Some(wait) = context.rate_limit.remaining(Instant::now()) { + return Err(RobotLBError::RateLimited(wait)); + } + tracing::info!("Starting service reconcilation"); let lb = LoadBalancer::try_from_svc(&svc, &context)?; @@ -526,9 +606,23 @@ async fn clear_ingress_status(svc_api: &kube::Api, svc: &Service) -> Ro /// Handle the error during reconcilation. #[allow(clippy::needless_pass_by_value)] -fn on_error(_: Arc, error: &RobotLBError, _context: Arc) -> Action { +fn on_error(svc: Arc, error: &RobotLBError, context: Arc) -> Action { + let service = format!("{}/{}", svc.namespace().unwrap_or_default(), svc.name_any()); + error_action(error, &context.rate_limit, Instant::now(), &service) +} + +fn error_action( + error: &RobotLBError, + rate_limit: &RateLimitGate, + now: Instant, + service: &str, +) -> Action { match error { RobotLBError::SkipService => Action::await_change(), + RobotLBError::RateLimited(wait) => Action::requeue(spread(*wait, service)), + error if error.is_rate_limited() => { + Action::requeue(spread(rate_limit.on_rate_limited(now), service)) + } _ => Action::requeue(Duration::from_secs(30)), } } @@ -536,8 +630,8 @@ fn on_error(_: Arc, error: &RobotLBError, _context: Arc #[cfg(test)] mod tests { use super::{ - collect_lb_services, consts, is_excluded_from_lb, is_lb_eligible_node, - is_local_traffic_policy, node_source, NodeSource, + collect_lb_services, consts, error_action, event_note, is_excluded_from_lb, + is_lb_eligible_node, is_local_traffic_policy, node_source, publishes_event, NodeSource, }; use k8s_openapi::{ api::core::v1::{ @@ -733,4 +827,85 @@ mod tests { }); assert_eq!(node_source(&svc, false), NodeSource::Annotation); } + + #[test] + fn an_event_note_is_redacted_and_bounded() { + let token = "0123456789abcdef"; + let error = crate::error::RobotLBError::HCloudError(format!( + "rejected token {} {}", + &token[..10], + "é".repeat(600) + )); + let note = event_note(&error, token); + assert!(!note.contains(&token[..10])); + assert!(note.contains("[REDACTED]")); + assert!(note.len() <= 1024); + } + + #[test] + fn a_gated_service_waits_out_the_pause() { + let gate = crate::rate_limit::RateLimitGate::default(); + let wait = std::time::Duration::from_secs(42); + let error = crate::error::RobotLBError::RateLimited(wait); + assert_eq!( + error_action(&error, &gate, std::time::Instant::now(), "shop/web"), + kube::runtime::controller::Action::requeue(crate::rate_limit::spread(wait, "shop/web")) + ); + } + + #[test] + fn a_rate_limited_call_closes_the_gate() { + let gate = crate::rate_limit::RateLimitGate::default(); + let now = std::time::Instant::now(); + let error = crate::error::RobotLBError::from(hcloud::apis::Error::< + hcloud::apis::load_balancers_api::AddTargetError, + >::ResponseError( + hcloud::apis::ResponseContent { + status: 429_u16.try_into().unwrap(), + content: String::new(), + entity: None, + }, + )); + assert_eq!( + error_action(&error, &gate, now, "shop/web"), + kube::runtime::controller::Action::requeue(crate::rate_limit::spread( + std::time::Duration::from_secs(60), + "shop/web" + )) + ); + assert!(gate.remaining(now).is_some()); + } + + #[test] + fn other_errors_retry_in_30_seconds() { + let gate = crate::rate_limit::RateLimitGate::default(); + let error = crate::error::RobotLBError::HCloudError("boom".to_string()); + assert_eq!( + error_action(&error, &gate, std::time::Instant::now(), "shop/web"), + kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(30)) + ); + } + + #[test] + fn every_failure_except_a_skip_publishes_an_event() { + use crate::error::RobotLBError; + assert!(!publishes_event(&RobotLBError::SkipService)); + // Jitter wakes the same service first after every pause, so the others only + // ever see the closed gate, and without an event of their own they go silent. + assert!(publishes_event(&RobotLBError::RateLimited( + std::time::Duration::from_secs(1) + ))); + assert!(publishes_event(&RobotLBError::from(hcloud::apis::Error::< + hcloud::apis::load_balancers_api::ListLoadBalancersError, + >::ResponseError( + hcloud::apis::ResponseContent { + status: 429_u16.try_into().unwrap(), + content: String::new(), + entity: None, + } + )))); + assert!(publishes_event(&RobotLBError::HCloudError( + "boom".to_string() + ))); + } } diff --git a/src/rate_limit.rs b/src/rate_limit.rs new file mode 100644 index 0000000..2792e08 --- /dev/null +++ b/src/rate_limit.rs @@ -0,0 +1,176 @@ +use std::{ + collections::hash_map::DefaultHasher, + hash::{Hash, Hasher}, + sync::Mutex, + time::{Duration, Instant}, +}; + +const FIRST_DELAY: Duration = Duration::from_secs(60); +const MAX_DOUBLINGS: u32 = 4; + +/// Pause shared by every service. Hetzner counts API requests per project, so once one +/// reconciliation is rate limited, any other call would only spend the budget being waited for. +/// +/// The generated `hcloud` client drops response headers, so `RateLimit-Reset` is not +/// available and the pause grows exponentially instead. +#[derive(Debug, Default)] +pub struct RateLimitGate { + state: Mutex, +} + +#[derive(Debug, Default)] +struct State { + closed_until: Option, + last_delay: Duration, + consecutive: u32, +} + +impl State { + fn remaining(&self, now: Instant) -> Option { + self.closed_until + .and_then(|until| until.checked_duration_since(now)) + .filter(|remaining| !remaining.is_zero()) + } +} + +impl RateLimitGate { + /// How long calls to the API must still wait, if they must. + pub fn remaining(&self, now: Instant) -> Option { + let state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + state.remaining(now) + } + + /// Record a rate-limited call and return how long to wait before the next one. + /// A call rejected while the gate is already closed was sent before it closed, + /// so it keeps the current pause instead of lengthening it. The pause starts short + /// again only after the API went without a 429 for as long as the last pause: + /// a reconciliation that succeeds may not have called the API at all. + pub fn on_rate_limited(&self, now: Instant) -> Duration { + let mut state = self + .state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + if let Some(remaining) = state.remaining(now) { + return remaining; + } + if let Some(reopened) = state.closed_until { + if now >= reopened + state.last_delay { + state.consecutive = 0; + } + } + let delay = FIRST_DELAY * 2_u32.pow(state.consecutive.min(MAX_DOUBLINGS)); + state.consecutive += 1; + state.closed_until = Some(now + delay); + state.last_delay = delay; + delay + } +} + +/// Push a wait out by up to a quarter, differently for each service, so services +/// paused together do not all call the API the moment the pause ends. +#[must_use] +pub fn spread(wait: Duration, service: &str) -> Duration { + let mut hasher = DefaultHasher::new(); + service.hash(&mut hasher); + let window = u64::try_from((wait / 4).as_millis()) + .unwrap_or(u64::MAX) + .max(1); + wait + Duration::from_millis(hasher.finish() % window) +} + +#[cfg(test)] +mod tests { + use super::{spread, RateLimitGate}; + use std::time::{Duration, Instant}; + + #[test] + fn an_untouched_gate_is_open() { + assert_eq!(RateLimitGate::default().remaining(Instant::now()), None); + } + + #[test] + fn a_rate_limit_closes_the_gate_for_the_returned_delay() { + let gate = RateLimitGate::default(); + let now = Instant::now(); + let delay = gate.on_rate_limited(now); + assert_eq!(delay, Duration::from_secs(60)); + assert_eq!(gate.remaining(now), Some(delay)); + assert_eq!(gate.remaining(now + delay), None); + } + + #[test] + fn repeated_rate_limits_double_the_delay_up_to_the_cap() { + let gate = RateLimitGate::default(); + let mut now = Instant::now(); + let mut delays = Vec::new(); + for _ in 0..8 { + let delay = gate.on_rate_limited(now); + delays.push(delay.as_secs()); + now += delay; + } + assert_eq!(delays, vec![60, 120, 240, 480, 960, 960, 960, 960]); + } + + #[test] + fn a_quiet_period_as_long_as_the_last_pause_resets_the_delay() { + let gate = RateLimitGate::default(); + let now = Instant::now(); + let first = gate.on_rate_limited(now); + let second = gate.on_rate_limited(now + first); + let reopened = now + first + second; + assert_eq!( + gate.on_rate_limited(reopened + second), + Duration::from_secs(60) + ); + } + + #[test] + fn a_shorter_quiet_period_keeps_doubling() { + let gate = RateLimitGate::default(); + let now = Instant::now(); + let first = gate.on_rate_limited(now); + let second = gate.on_rate_limited(now + first); + let reopened = now + first + second; + assert_eq!( + gate.on_rate_limited(reopened + second.checked_sub(Duration::from_secs(1)).unwrap()), + Duration::from_secs(240) + ); + } + + #[test] + fn calls_already_in_flight_do_not_lengthen_the_pause() { + let gate = RateLimitGate::default(); + let now = Instant::now(); + gate.on_rate_limited(now); + let later = now + Duration::from_secs(10); + assert_eq!(gate.on_rate_limited(later), Duration::from_secs(50)); + assert_eq!( + gate.on_rate_limited(later + Duration::from_secs(50)), + Duration::from_secs(120) + ); + } + + #[test] + fn spread_adds_up_to_a_quarter_of_the_pause() { + let wait = Duration::from_secs(60); + for name in ["web", "api", "dns", "ingress-nginx-controller"] { + let spread_wait = spread(wait, name); + assert!(spread_wait >= wait); + assert!(spread_wait < wait + wait / 4); + } + } + + #[test] + fn spread_is_stable_per_service_and_differs_between_services() { + let wait = Duration::from_secs(60); + assert_eq!(spread(wait, "web"), spread(wait, "web")); + let wakeups = ["a", "b", "c", "d", "e", "f"] + .map(|name| spread(wait, name)) + .into_iter() + .collect::>(); + assert!(wakeups.len() > 1); + } +}