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
128 changes: 128 additions & 0 deletions src/backoff.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,128 @@
use std::{
collections::HashMap,
sync::Mutex,
time::{Duration, Instant},
};

const FIRST_DELAY: Duration = Duration::from_secs(5);
const MAX_DELAY: Duration = Duration::from_secs(300);
/// A service deleted before it got the robotlb finalizer is not reconciled again to be
/// forgotten on success, so an entry that has not failed for this long is dropped. It is
/// well above `MAX_DELAY` plus one rate limit pause; a service held at the rate limit gate
/// for longer starts short again.
const FORGET_AFTER: Duration = Duration::from_secs(3600);

/// Retry delay of each service whose reconciliation keeps failing.
///
/// It doubles from `FIRST_DELAY` up to `MAX_DELAY` like the upstream cloud-provider
/// service controller:
/// <https://github.com/kubernetes/cloud-provider/blob/master/controllers/service/controller.go>
#[derive(Debug, Default)]
pub struct ErrorBackoff {
failures: Mutex<HashMap<String, Failures>>,
}

#[derive(Debug)]
struct Failures {
count: u32,
last: Instant,
}

impl ErrorBackoff {
/// Record a failed reconciliation of `service` and return how long to wait before the next one.
pub fn on_failure(&self, service: &str, now: Instant) -> Duration {
let mut failures = self
.failures
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
failures.retain(|_, failure| now.saturating_duration_since(failure.last) <= FORGET_AFTER);
let failure = failures.entry(service.to_string()).or_insert(Failures {
count: 0,
last: now,
});
let delay = FIRST_DELAY
.saturating_mul(2_u32.saturating_pow(failure.count))
.min(MAX_DELAY);
failure.count = failure.count.saturating_add(1);
failure.last = now;
drop(failures);
delay
}

/// Start the delay of `service` short again.
pub fn forget(&self, service: &str) {
self.failures
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.remove(service);
}

#[cfg(test)]
fn tracked(&self) -> usize {
self.failures
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.len()
}
}

#[cfg(test)]
mod tests {
use super::{ErrorBackoff, FORGET_AFTER};
use std::time::{Duration, Instant};

#[test]
fn repeated_failures_double_the_delay_up_to_the_cap() {
let backoff = ErrorBackoff::default();
let mut now = Instant::now();
let mut delays = Vec::new();
for _ in 0..9 {
let delay = backoff.on_failure("shop/web", now);
delays.push(delay.as_secs());
now += delay;
}
assert_eq!(delays, vec![5, 10, 20, 40, 80, 160, 300, 300, 300]);
}

#[test]
fn services_back_off_independently() {
let backoff = ErrorBackoff::default();
let now = Instant::now();
backoff.on_failure("shop/web", now);
backoff.on_failure("shop/web", now);
assert_eq!(backoff.on_failure("shop/api", now), Duration::from_secs(5));
assert_eq!(backoff.on_failure("shop/web", now), Duration::from_secs(20));
}

#[test]
fn a_forgotten_service_starts_short_again() {
let backoff = ErrorBackoff::default();
let now = Instant::now();
backoff.on_failure("shop/web", now);
backoff.on_failure("shop/web", now);
backoff.forget("shop/web");
assert_eq!(backoff.on_failure("shop/web", now), Duration::from_secs(5));
}

#[test]
fn a_long_failing_service_is_not_forgotten() {
let backoff = ErrorBackoff::default();
let mut now = Instant::now();
for _ in 0..20 {
now += backoff.on_failure("shop/web", now);
}
assert_eq!(
backoff.on_failure("shop/web", now),
Duration::from_secs(300)
);
}

#[test]
fn services_that_stopped_failing_are_dropped() {
let backoff = ErrorBackoff::default();
let now = Instant::now();
backoff.on_failure("shop/deleted", now);
backoff.on_failure("shop/web", now + FORGET_AFTER + Duration::from_secs(1));
assert_eq!(backoff.tracked(), 1);
}
}
93 changes: 85 additions & 8 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
)
]

use backoff::ErrorBackoff;
use clap::Parser;
use config::OperatorConfig;
use error::{redact, RobotLBError, RobotLBResult};
Expand Down Expand Up @@ -47,6 +48,7 @@ use std::{
time::{Duration, Instant},
};

pub mod backoff;
pub mod config;
pub mod consts;
pub mod error;
Expand Down Expand Up @@ -124,6 +126,7 @@ pub struct CurrentContext {
pub config: OperatorConfig,
pub hcloud_config: HCloudConfig,
pub rate_limit: Arc<RateLimitGate>,
pub error_backoff: Arc<ErrorBackoff>,
}
impl CurrentContext {
#[must_use]
Expand All @@ -133,6 +136,7 @@ impl CurrentContext {
config,
hcloud_config,
rate_limit: Arc::default(),
error_backoff: Arc::default(),
}
}
}
Expand All @@ -147,6 +151,9 @@ pub async fn reconcile_service(
context: Arc<CurrentContext>,
) -> RobotLBResult<Action> {
let result = sync_service(svc.clone(), context.clone()).await;
if ends_backoff(&result) {
context.error_backoff.forget(&service_key(&svc));
}
if let Err(error) = &result {
if publishes_event(error) {
report_failure(&svc, &context, error).await;
Expand All @@ -155,6 +162,16 @@ pub async fn reconcile_service(
result
}

/// A skipped service may have failed back when it was still a robotlb load balancer.
/// A service waiting at the rate limit gate has not got any closer to succeeding.
const fn ends_backoff(result: &RobotLBResult<Action>) -> bool {
matches!(result, Ok(_) | Err(RobotLBError::SkipService))
}

fn service_key(svc: &Service) -> String {
format!("{}/{}", svc.namespace().unwrap_or_default(), svc.name_any())
}

/// 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 {
Expand Down Expand Up @@ -607,13 +624,19 @@ async fn clear_ingress_status(svc_api: &kube::Api<Service>, svc: &Service) -> Ro
/// Handle the error during reconcilation.
#[allow(clippy::needless_pass_by_value)]
fn on_error(svc: Arc<Service>, error: &RobotLBError, context: Arc<CurrentContext>) -> Action {
let service = format!("{}/{}", svc.namespace().unwrap_or_default(), svc.name_any());
error_action(error, &context.rate_limit, Instant::now(), &service)
error_action(
error,
&context.rate_limit,
&context.error_backoff,
Instant::now(),
&service_key(&svc),
)
}

fn error_action(
error: &RobotLBError,
rate_limit: &RateLimitGate,
backoff: &ErrorBackoff,
now: Instant,
service: &str,
) -> Action {
Expand All @@ -623,7 +646,7 @@ fn error_action(
error if error.is_rate_limited() => {
Action::requeue(spread(rate_limit.on_rate_limited(now), service))
}
_ => Action::requeue(Duration::from_secs(30)),
_ => Action::requeue(backoff.on_failure(service, now)),
}
}

Expand Down Expand Up @@ -848,7 +871,13 @@ 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(), "shop/web"),
error_action(
&error,
&gate,
&crate::backoff::ErrorBackoff::default(),
std::time::Instant::now(),
"shop/web"
),
kube::runtime::controller::Action::requeue(crate::rate_limit::spread(wait, "shop/web"))
);
}
Expand All @@ -867,7 +896,13 @@ mod tests {
},
));
assert_eq!(
error_action(&error, &gate, now, "shop/web"),
error_action(
&error,
&gate,
&crate::backoff::ErrorBackoff::default(),
now,
"shop/web"
),
kube::runtime::controller::Action::requeue(crate::rate_limit::spread(
std::time::Duration::from_secs(60),
"shop/web"
Expand All @@ -877,15 +912,57 @@ mod tests {
}

#[test]
fn other_errors_retry_in_30_seconds() {
fn other_errors_back_off_per_service() {
let gate = crate::rate_limit::RateLimitGate::default();
let backoff = crate::backoff::ErrorBackoff::default();
let now = std::time::Instant::now();
let error = crate::error::RobotLBError::HCloudError("boom".to_string());
let retry = |service| error_action(&error, &gate, &backoff, now, service);
let after =
|secs| kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(secs));
assert_eq!(retry("shop/web"), after(5));
assert_eq!(retry("shop/web"), after(10));
assert_eq!(retry("shop/api"), after(5));
}

#[test]
fn rate_limits_do_not_lengthen_the_error_backoff() {
let gate = crate::rate_limit::RateLimitGate::default();
let backoff = crate::backoff::ErrorBackoff::default();
let now = std::time::Instant::now();
let gated = crate::error::RobotLBError::RateLimited(std::time::Duration::from_secs(42));
error_action(&gated, &gate, &backoff, now, "shop/web");
let rejected = 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,
},
));
error_action(&rejected, &gate, &backoff, now, "shop/web");
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))
error_action(&error, &gate, &backoff, now, "shop/web"),
kube::runtime::controller::Action::requeue(std::time::Duration::from_secs(5))
);
}

#[test]
fn success_and_skips_end_the_backoff_but_a_closed_gate_does_not() {
use crate::error::RobotLBError;
use kube::runtime::controller::Action;
assert!(super::ends_backoff(&Ok(Action::await_change())));
assert!(super::ends_backoff(&Err(RobotLBError::SkipService)));
assert!(!super::ends_backoff(&Err(RobotLBError::RateLimited(
std::time::Duration::from_secs(1)
))));
assert!(!super::ends_backoff(&Err(RobotLBError::HCloudError(
"boom".to_string()
))));
}

#[test]
fn every_failure_except_a_skip_publishes_an_event() {
use crate::error::RobotLBError;
Expand Down
Loading