From 0b8d7948079d84f4c035590c56fd38ffdd850f53 Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 00:50:26 +0300 Subject: [PATCH 1/4] fix(hcloud): back off on rate limits and report failures on the service Once the Hetzner Cloud API budget ran out, every service was retried every 30 seconds, and each retry spent the budget being waited for. The budget is per project, so after a 429 all services now pause together, starting at one minute and doubling up to 16 minutes while the limit keeps being hit. Adding targets stops at the first 429 instead of calling the API for every remaining node. The generated client drops response headers, so RateLimit-Reset cannot be used. A failed reconciliation now leaves a SyncLoadBalancerFailed warning event on the service, so kubectl describe shows the last error. Errors are logged on one line with the status and the message Hetzner reported, and the API token, which Hetzner quotes in some messages, is redacted from logs and events. The chart role gains permission to create events. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- helm/values.yaml | 3 + src/error.rs | 186 +++++++++++++++++++++++++++++++++++++--- src/lb.rs | 17 ++-- src/main.rs | 213 ++++++++++++++++++++++++++++++++++++++++------ src/rate_limit.rs | 141 ++++++++++++++++++++++++++++++ 5 files changed, 512 insertions(+), 48 deletions(-) create mode 100644 src/rate_limit.rs diff --git a/helm/values.yaml b/helm/values.yaml index 2524ad9..85c559a 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -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..d36c3ec 100644 --- a/src/error.rs +++ b/src/error.rs @@ -22,65 +22,223 @@ pub enum RobotLBError { UnknownLBAlgorithm, #[error("Cannot get target nodes, because the service has no selector")] ServiceWithoutSelector, + #[error("Hetzner Cloud API rate limit reached, next attempt in {}s", .0.as_secs())] + 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_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..0f6e964 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::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::debug!("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,63 @@ 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, and a gated one has not +/// called the API at all, so neither has a failure to show. +const fn publishes_event(error: &RobotLBError) -> bool { + !matches!( + error, + RobotLBError::SkipService | RobotLBError::RateLimited(_) + ) +} + +/// 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 +225,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 +609,15 @@ 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(_: Arc, error: &RobotLBError, context: Arc) -> Action { + error_action(error, &context.rate_limit, Instant::now()) +} + +fn error_action(error: &RobotLBError, rate_limit: &RateLimitGate, now: Instant) -> Action { match error { RobotLBError::SkipService => Action::await_change(), + RobotLBError::RateLimited(wait) => Action::requeue(*wait), + error if error.is_rate_limited() => Action::requeue(rate_limit.on_rate_limited(now)), _ => Action::requeue(Duration::from_secs(30)), } } @@ -536,8 +625,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 +822,80 @@ 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()), + kube::runtime::controller::Action::requeue(wait) + ); + } + + #[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), + kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(60)) + ); + 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()), + kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(30)) + ); + } + + #[test] + fn only_real_failures_publish_an_event() { + use crate::error::RobotLBError; + assert!(!publishes_event(&RobotLBError::SkipService)); + 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..bfa2e80 --- /dev/null +++ b/src/rate_limit.rs @@ -0,0 +1,141 @@ +use std::{ + 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 + } +} + +#[cfg(test)] +mod tests { + use super::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) + ); + } +} From 851e4a3d4706e5463e20bc8a6fa0932cf33ba7ff Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 12:18:40 +0300 Subject: [PATCH 2/4] fix(hcloud): spread wakeups after a rate limit pause Services paused by the rate limit gate were all requeued for the moment it reopened, so they hit the API in one burst while the budget was still nearly empty. Each wait now grows by up to a quarter, derived from the service name, so services wake up spread over that window. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- src/main.rs | 33 ++++++++++++++++++++++----------- src/rate_limit.rs | 37 ++++++++++++++++++++++++++++++++++++- 2 files changed, 58 insertions(+), 12 deletions(-) diff --git a/src/main.rs b/src/main.rs index 0f6e964..6ec1ef9 100644 --- a/src/main.rs +++ b/src/main.rs @@ -39,7 +39,7 @@ use kube::{ }; use label_filter::LabelFilter; use lb::{LBService, LoadBalancer}; -use rate_limit::RateLimitGate; +use rate_limit::{spread, RateLimitGate}; use std::{ collections::HashSet, str::FromStr, @@ -609,15 +609,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 { - error_action(error, &context.rate_limit, Instant::now()) +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) -> Action { +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(*wait), - error if error.is_rate_limited() => Action::requeue(rate_limit.on_rate_limited(now)), + 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)), } } @@ -843,8 +851,8 @@ mod tests { 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()), - kube::runtime::controller::Action::requeue(wait) + error_action(&error, &gate, std::time::Instant::now(), "shop/web"), + kube::runtime::controller::Action::requeue(crate::rate_limit::spread(wait, "shop/web")) ); } @@ -862,8 +870,11 @@ mod tests { }, )); assert_eq!( - error_action(&error, &gate, now), - kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(60)) + 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()); } @@ -873,7 +884,7 @@ mod tests { 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()), + error_action(&error, &gate, std::time::Instant::now(), "shop/web"), kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(30)) ); } diff --git a/src/rate_limit.rs b/src/rate_limit.rs index bfa2e80..2792e08 100644 --- a/src/rate_limit.rs +++ b/src/rate_limit.rs @@ -1,4 +1,6 @@ use std::{ + collections::hash_map::DefaultHasher, + hash::{Hash, Hasher}, sync::Mutex, time::{Duration, Instant}, }; @@ -67,9 +69,21 @@ impl RateLimitGate { } } +/// 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::RateLimitGate; + use super::{spread, RateLimitGate}; use std::time::{Duration, Instant}; #[test] @@ -138,4 +152,25 @@ mod tests { 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); + } } From d39dc1cd077e37bef9e74cff0041f929207b8c5e Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 12:41:05 +0300 Subject: [PATCH 3/4] fix(hcloud): report services waiting at the rate limit gate With wakeups spread by service name, the same service wakes first after every pause, hits the limit and closes the gate again. Every other service only ever saw the closed gate, which published no event and logged at debug level, so after the event TTL they showed nothing at all. A service that wakes up to a closed gate now gets an event and an info line. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- src/error.rs | 8 +++++++- src/main.rs | 17 ++++++++--------- 2 files changed, 15 insertions(+), 10 deletions(-) diff --git a/src/error.rs b/src/error.rs index d36c3ec..ec91801 100644 --- a/src/error.rs +++ b/src/error.rs @@ -22,7 +22,7 @@ pub enum RobotLBError { UnknownLBAlgorithm, #[error("Cannot get target nodes, because the service has no selector")] ServiceWithoutSelector, - #[error("Hetzner Cloud API rate limit reached, next attempt in {}s", .0.as_secs())] + #[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 @@ -223,6 +223,12 @@ mod tests { 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, "")); diff --git a/src/main.rs b/src/main.rs index 6ec1ef9..edcfa73 100644 --- a/src/main.rs +++ b/src/main.rs @@ -100,7 +100,7 @@ async fn main() -> RobotLBResult<()> { error @ RobotLBError::RateLimited(_), service, )) => { - tracing::debug!("Service {service}: {error}"); + tracing::info!("Service {service}: {error}"); } Err(controller::Error::ReconcilerFailed(error, service)) => { tracing::error!( @@ -155,13 +155,10 @@ pub async fn reconcile_service( result } -/// Skipped services are every service robotlb does not own, and a gated one has not -/// called the API at all, so neither has a failure to show. +/// 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 | RobotLBError::RateLimited(_) - ) + !matches!(error, RobotLBError::SkipService) } /// Put the error on the service as a warning event, where `kubectl describe` shows it. @@ -890,10 +887,12 @@ mod tests { } #[test] - fn only_real_failures_publish_an_event() { + fn every_failure_except_a_skip_publishes_an_event() { use crate::error::RobotLBError; assert!(!publishes_event(&RobotLBError::SkipService)); - assert!(!publishes_event(&RobotLBError::RateLimited( + // 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::< From 785e5bf1ee8ad4452889170aba4d2773504a9905 Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 12:41:06 +0300 Subject: [PATCH 4/4] docs(helm): describe the default service account permissions The comment claimed the default grants all permissions, while it lists only what robotlb needs. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- helm/values.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/helm/values.yaml b/helm/values.yaml index 85c559a..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]