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)