Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 77 additions & 0 deletions src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,29 @@ pub fn describe<T>(error: &hcloud::apis::Error<T>) -> 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<T>(error: &hcloud::apis::Error<T>) -> bool {
let hcloud::apis::Error::ResponseError(response) = error else {
return false;
};
let code = k8s_openapi::serde_json::from_str::<k8s_openapi::serde_json::Value>(
&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]
Expand Down Expand Up @@ -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,
"<html></html>"
)));
}

#[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<ListLoadBalancersError> =
Error::Serde(k8s_openapi::serde_json::from_str::<()>("x").unwrap_err());
assert!(!super::is_temporary_rejection(&error));
}
}
142 changes: 130 additions & 12 deletions src/lb.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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 {
Expand Down Expand Up @@ -703,6 +719,33 @@ impl From<LBAlgorithm> 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<T, E, F, Fut>(
retries: u32,
unit: std::time::Duration,
mut call: F,
) -> Result<T, hcloud::apis::Error<E>>
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = Result<T, hcloud::apis::Error<E>>>,
{
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;
Expand All @@ -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<AddTargetError> {
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<AddTargetError>>, 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)
Expand Down
Loading