From 128fd311e7dd2a5c8e6b3d4f76f2b9b5d1b653c0 Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 19:07:17 +0300 Subject: [PATCH] fix(lb): retry targets rejected because the balancer is busy Hetzner answers 423 locked while an action is still running on the balancer, and the operator skipped the target on that answer as if it were permanently invalid. A node added right after another change could stay out of the balancer until the next reconcile. Tell temporary rejections (locked, conflict, robot_unavailable, 5xx) from permanent ones and retry the former up to twice, after 1s and 2s. A target that still fails is skipped as before, and a 429 is still returned at once for the rate limit gate. Once one target has used up its retries, the rest of the run does not retry: the lock is on the whole balancer, and waiting per target would stall the service for minutes on a large balancer. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- src/error.rs | 77 ++++++++++++++++++++++++++++ src/lb.rs | 142 ++++++++++++++++++++++++++++++++++++++++++++++----- 2 files changed, 207 insertions(+), 12 deletions(-) diff --git a/src/error.rs b/src/error.rs index ec91801..15616df 100644 --- a/src/error.rs +++ b/src/error.rs @@ -152,6 +152,29 @@ pub fn describe(error: &hcloud::apis::Error) -> String { ) } +/// Whether Hetzner refused the call for a reason that may pass on its own. +/// +/// That is a balancer locked by a running action (error code `locked`, HTTP 423), a +/// resource changed during the request (`conflict`, HTTP 409), Robot being briefly +/// unavailable (`robot_unavailable`), or a failure on the API side (5xx). +/// A 429 is not temporary here: the rate limit gate owns it. +#[must_use] +pub fn is_temporary_rejection(error: &hcloud::apis::Error) -> bool { + let hcloud::apis::Error::ResponseError(response) = error else { + return false; + }; + let code = k8s_openapi::serde_json::from_str::( + &response.content, + ) + .map(|body| { + body.pointer("/error/code") + .and_then(|code| code.as_str()) + .map(str::to_owned) + }); + let retryable = matches!(code, Ok(Some(code)) if matches!(code.as_str(), "locked" | "conflict" | "robot_unavailable")); + retryable || response.status.is_server_error() +} + /// 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] @@ -247,4 +270,58 @@ mod tests { assert!(!RobotLBError::from(response_error(500, "")).is_rate_limited()); assert!(!RobotLBError::SkipService.is_rate_limited()); } + + #[test] + fn a_locked_balancer_is_temporary() { + let body = r#"{"error": {"code": "locked", "message": "item is locked"}}"#; + assert!(super::is_temporary_rejection(&response_error(423, body))); + assert!(super::is_temporary_rejection(&response_error(409, body))); + } + + #[test] + fn a_conflicting_change_is_temporary() { + let body = r#"{"error": {"code": "conflict", "message": "please retry"}}"#; + assert!(super::is_temporary_rejection(&response_error(409, body))); + } + + #[test] + fn an_unavailable_robot_is_temporary() { + let body = r#"{"error": {"code": "robot_unavailable", "message": "retry later"}}"#; + assert!(super::is_temporary_rejection(&response_error(422, body))); + } + + #[test] + fn a_protected_balancer_is_permanent() { + let body = r#"{"error": {"code": "protected", "message": "protected"}}"#; + assert!(!super::is_temporary_rejection(&response_error(423, body))); + } + + #[test] + fn a_server_error_is_temporary() { + assert!(super::is_temporary_rejection(&response_error(500, ""))); + assert!(super::is_temporary_rejection(&response_error( + 503, + "" + ))); + } + + #[test] + fn a_rate_limit_is_not_temporary() { + let body = r#"{"error": {"code": "rate_limit_exceeded", "message": "slow down"}}"#; + assert!(!super::is_temporary_rejection(&response_error(429, body))); + } + + #[test] + fn a_target_outside_the_subnet_is_permanent() { + let body = r#"{"error": {"code": "invalid_input", "message": "not in subnet"}}"#; + assert!(!super::is_temporary_rejection(&response_error(422, body))); + assert!(!super::is_temporary_rejection(&response_error(404, ""))); + } + + #[test] + fn a_transport_failure_is_not_temporary() { + let error: Error = + Error::Serde(k8s_openapi::serde_json::from_str::<()>("x").unwrap_err()); + assert!(!super::is_temporary_rejection(&error)); + } } diff --git a/src/lb.rs b/src/lb.rs index edd1478..70ddd1a 100644 --- a/src/lb.rs +++ b/src/lb.rs @@ -26,6 +26,10 @@ use crate::{ CurrentContext, }; +/// Retries after the first attempt, sleeping 1 unit, 2 units, ... between them. +const ADD_TARGET_RETRIES: u32 = 2; +const ADD_TARGET_RETRY_UNIT: std::time::Duration = std::time::Duration::from_secs(1); + #[derive(Debug)] pub struct LBService { pub listen_port: i32, @@ -340,6 +344,7 @@ impl LoadBalancer { let mut live = 0_usize; let mut last_error = None; + let mut retries = ADD_TARGET_RETRIES; for ip in &planned { if hcloud_balancer .targets @@ -350,19 +355,30 @@ impl LoadBalancer { continue; } tracing::info!("Adding target {}", ip); - let added = hcloud::apis::load_balancers_api::add_target( - &self.hcloud_config, - AddTargetParams { - id: hcloud_balancer.id, - body: Some(LoadBalancerAddTarget { - ip: Some(Box::new(hcloud::models::LoadBalancerTargetIp { - ip: (*ip).to_string(), - })), - ..Default::default() - }), - }, - ) + let added = retry_temporary(retries, ADD_TARGET_RETRY_UNIT, || { + hcloud::apis::load_balancers_api::add_target( + &self.hcloud_config, + AddTargetParams { + id: hcloud_balancer.id, + body: Some(LoadBalancerAddTarget { + ip: Some(Box::new(hcloud::models::LoadBalancerTargetIp { + ip: (*ip).to_string(), + })), + ..Default::default() + }), + }, + ) + }) .await; + // A lock is on the whole balancer and a failing API stays failing for a while, + // so once the retries ran out for one target the remaining ones would only + // wait the same time for nothing. + if added + .as_ref() + .is_err_and(crate::error::is_temporary_rejection) + { + retries = 0; + } // Hetzner rejects IPs outside the vSwitch subnet of the attached network, // which must not keep the remaining nodes out of the load balancer. match added { @@ -703,6 +719,33 @@ impl From for LoadBalancerAlgorithm { } } +/// Call `call`, and call it again after `unit`, `2 * unit`, ... while Hetzner rejects it +/// for a temporary reason, up to `retries` more times. Returns the last outcome. +async fn retry_temporary( + retries: u32, + unit: std::time::Duration, + mut call: F, +) -> Result> +where + F: FnMut() -> Fut, + Fut: std::future::Future>>, +{ + let mut attempt = 0_u32; + loop { + match call().await { + Err(error) if attempt < retries && crate::error::is_temporary_rejection(&error) => { + attempt += 1; + tracing::debug!( + "Rejected temporarily, retry {attempt}: {}", + crate::error::describe(&error) + ); + tokio::time::sleep(unit * attempt).await; + } + outcome => return outcome, + } + } +} + #[cfg(test)] mod tests { use super::plan_targets; @@ -720,6 +763,81 @@ mod tests { ); } + use super::retry_temporary; + use hcloud::apis::{load_balancers_api::AddTargetError, Error, ResponseContent}; + use std::{ + sync::atomic::{AtomicU32, Ordering}, + time::Duration, + }; + + fn rejected(status: u16) -> Error { + Error::ResponseError(ResponseContent { + status: status.try_into().unwrap(), + content: if status == 423 { + r#"{"error": {"code": "locked", "message": "item is locked"}}"#.to_string() + } else { + String::new() + }, + entity: None, + }) + } + + /// Answers with `statuses` in order, then with success; returns the outcome and the call count. + async fn run(retries: u32, statuses: &[u16]) -> (Result<(), Error>, u32) { + let calls = AtomicU32::new(0); + let result = retry_temporary(retries, Duration::ZERO, || { + let call = calls.fetch_add(1, Ordering::Relaxed); + let answer = statuses + .get(call as usize) + .map_or(Ok(()), |s| Err(rejected(*s))); + async move { answer } + }) + .await; + (result, calls.load(Ordering::Relaxed)) + } + + #[tokio::test] + async fn a_locked_balancer_is_retried_until_it_accepts() { + let (result, calls) = run(2, &[423, 423]).await; + assert!(result.is_ok()); + assert_eq!(calls, 3); + } + + #[tokio::test] + async fn a_balancer_that_stays_locked_gets_the_first_call_and_the_retries_only() { + let (result, calls) = run(2, &[423; 10]).await; + assert!(result.is_err()); + assert_eq!(calls, 3); + } + + #[tokio::test] + async fn no_retries_means_a_single_call() { + let (result, calls) = run(0, &[423; 10]).await; + assert!(result.is_err()); + assert_eq!(calls, 1); + } + + #[tokio::test] + async fn a_permanent_rejection_is_not_retried() { + let (result, calls) = run(2, &[422; 10]).await; + assert!(result.is_err()); + assert_eq!(calls, 1); + } + + #[tokio::test] + async fn a_rate_limit_is_not_retried() { + let (result, calls) = run(2, &[429; 10]).await; + assert!(result.is_err()); + assert_eq!(calls, 1); + } + + #[tokio::test] + async fn a_success_is_not_repeated() { + let (result, calls) = run(2, &[]).await; + assert!(result.is_ok()); + assert_eq!(calls, 1); + } + #[test] fn targets_beyond_the_balancer_limit_are_dropped() { let desired = (1..=30)