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
4 changes: 3 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,9 @@ After the chart is installed, you should be able to create `LoadBalancer` servic

The operator listens to the Kubernetes API for services of type `LoadBalancer` and creates Hetzner load balancers that point to nodes based on `node-ip`.

A balancer is updated when its service changes, when a node is added, removed, relabelled, cordoned or changes readiness or addresses, and, for services with `externalTrafficPolicy: Local` while `ROBOTLB_DYNAMIC_NODE_SELECTOR` is on, when the nodes their endpoints serve traffic from change; a deleted endpoint slice of such a service rechecks every balancer. Apart from that robotlb checks each balancer every `ROBOTLB_RESYNC_INTERVAL` seconds, 300 by default, which bounds how long a change made to the balancer in Hetzner survives. Each check costs one or two Hetzner API requests per service. When Hetzner refuses a target for a reason other than the rate limit, the service is checked again within 30 seconds instead, so a node it refuses for good, such as one outside the vSwitch subnet, keeps its service on that 30-second cycle.
robotlb handles every `LoadBalancer` service that has no `spec.loadBalancerClass` or has it set to `robotlb`, and ignores services of any other class. Another LoadBalancer controller that also takes services without a class, such as MetalLB started without `--lb-class`, handles the same services: the two keep overwriting each other's `status.loadBalancer`, and every status write starts another reconcile. Set `loadBalancerClass: robotlb` on the services robotlb should handle to keep other controllers away from them. Kubernetes lets you set the field only when the service is created or its type is changed to `LoadBalancer`. Recreating a service deletes its Hetzner balancer, so the service gets a new balancer with a new public IP. Set the class when the service is created.

A balancer is updated when its service changes, when a node is added, removed, relabelled, cordoned or changes readiness or addresses, and, for services with `externalTrafficPolicy: Local` while `ROBOTLB_DYNAMIC_NODE_SELECTOR` is on, when the nodes their endpoints serve traffic from change; a deleted endpoint slice of such a service rechecks every balancer. Apart from that robotlb checks each balancer every `ROBOTLB_RESYNC_INTERVAL` seconds, 300 by default, which bounds how long a change made to the balancer in Hetzner survives. Each check costs one or two Hetzner API requests per service. A reconcile that runs before the next check is due, and finds the spec, the `robotlb/` annotations, the targets and the ports as the last successful reconcile left them, makes no Hetzner requests and does not write the status, so a change made in Hetzner, or a status another controller cleared, stays until that check. When Hetzner refuses a target for a reason other than the rate limit, the service is checked again within 30 seconds instead, so a node it refuses for good, such as one outside the vSwitch subnet, keeps its service on that 30-second cycle.

Target nodes are selected according to the service's `externalTrafficPolicy`:

Expand Down
29 changes: 25 additions & 4 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ pub mod error;
pub mod finalizers;
pub mod label_filter;
pub mod lb;
pub mod memo;
pub mod rate_limit;
pub mod triggers;

Expand Down Expand Up @@ -222,6 +223,7 @@ fn cluster_changes(
#[derive(Clone)]
pub struct CurrentContext {
pub client: kube::Client,
pub memo: Arc<memo::ReconcileMemo>,
pub config: OperatorConfig,
pub hcloud_config: HCloudConfig,
pub rate_limit: Arc<RateLimitGate>,
Expand All @@ -231,6 +233,7 @@ impl CurrentContext {
pub fn new(client: kube::Client, config: OperatorConfig, hcloud_config: HCloudConfig) -> Self {
Self {
client,
memo: Arc::default(),
config,
hcloud_config,
rate_limit: Arc::default(),
Expand All @@ -253,6 +256,7 @@ pub async fn reconcile_service(
report_failure(&svc, &context, error).await;
}
}
context.memo.settle(&svc, result.is_ok());
result
}

Expand Down Expand Up @@ -568,11 +572,25 @@ pub async fn reconcile_load_balancer(
if lb.services.is_empty() {
tracing::warn!("Service has no port that can be exposed. Skipping the load balancer.");
clear_ingress_status(&svc_api, &svc).await?;
// The cleared status must come back once the ports are restored.
if let Some(uid) = svc.uid() {
context.memo.forget(&uid);
}
return Ok(Action::requeue(Duration::from_secs(
context.config.resync_interval,
)));
}

let uid = svc.uid();
let fingerprint = memo::fingerprint(&svc, &lb.targets, &lb.services);
if let Some(wait) = uid
.as_deref()
.and_then(|uid| context.memo.unchanged(uid, fingerprint, Instant::now()))
{
tracing::debug!("Nothing robotlb acts on changed since the last reconcile. Skipping...");
return Ok(Action::requeue(wait));
}

let (hcloud_lb, targets_missing) = lb.reconcile().await?;

let mut ingress = vec![];
Expand Down Expand Up @@ -614,10 +632,13 @@ pub async fn reconcile_load_balancer(
.await?;
}

Ok(Action::requeue(success_requeue(
targets_missing,
context.config.resync_interval,
)))
let requeue = success_requeue(targets_missing, context.config.resync_interval);
if let Some(uid) = &uid {
context
.memo
.record(uid, fingerprint, Instant::now(), requeue);
}
Ok(Action::requeue(requeue))
}

/// A target Hetzner rejected, for a reason other than the rate limit, may be accepted
Expand Down
299 changes: 299 additions & 0 deletions src/memo.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,299 @@
use k8s_openapi::api::core::v1::Service;
use kube::{Resource, ResourceExt};
use std::{
collections::{BTreeMap, BTreeSet, HashMap},
hash::{BuildHasher, DefaultHasher, Hash, Hasher},
sync::{Mutex, PoisonError},
time::{Duration, Instant},
};

/// What each service looked like after its last successful reconcile, so that a
/// watch event that changed nothing robotlb acts on costs no Hetzner requests.
#[derive(Default)]
pub struct ReconcileMemo {
entries: Mutex<HashMap<String, (u64, Instant)>>,
}

impl ReconcileMemo {
/// How long the recorded state stays valid, when the service is still in it.
pub fn unchanged(&self, uid: &str, fingerprint: u64, now: Instant) -> Option<Duration> {
let (recorded, until) = *self
.entries
.lock()
.unwrap_or_else(PoisonError::into_inner)
.get(uid)?;
if recorded != fingerprint {
return None;
}
Some(until.saturating_duration_since(now)).filter(|left| !left.is_zero())
}

/// Keep the state for `valid_for`, the time until the reconcile is due anyway.
pub fn record(&self, uid: &str, fingerprint: u64, now: Instant, valid_for: Duration) {
self.entries
.lock()
.unwrap_or_else(PoisonError::into_inner)
.insert(uid.to_string(), (fingerprint, now + valid_for));
}

/// Drop the entry of a service after a failed reconcile, or once robotlb no
/// longer serves it.
pub fn settle(&self, svc: &Service, succeeded: bool) {
let serves =
svc.meta().deletion_timestamp.is_none() && crate::is_robotlb_load_balancer(svc);
if !(succeeded && serves) {
if let Some(uid) = svc.uid() {
self.forget(&uid);
}
}
}

pub fn forget(&self, uid: &str) {
self.entries
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(uid);
}
}

/// A hash of everything a reconcile acts on: the spec, the robotlb annotations, and
/// the targets and ports computed from them.
pub fn fingerprint<S: BuildHasher>(
svc: &Service,
targets: &[String],
services: &HashMap<i32, i32, S>,
) -> u64 {
// `DefaultHasher::new` uses fixed keys, unlike `RandomState`, so the same input
// hashes the same on every call within the process.
let mut hasher = DefaultHasher::new();
k8s_openapi::serde_json::to_string(&svc.spec)
.unwrap_or_default()
.hash(&mut hasher);
for annotation in svc
.annotations()
.iter()
.filter(|(key, _)| key.starts_with("robotlb/"))
{
annotation.hash(&mut hasher);
}
targets.iter().collect::<BTreeSet<_>>().hash(&mut hasher);
services
.iter()
.collect::<BTreeMap<_, _>>()
.hash(&mut hasher);
hasher.finish()
}

#[cfg(test)]
mod tests {
use super::{fingerprint, ReconcileMemo};
use crate::consts;
use k8s_openapi::{
api::core::v1::{
LoadBalancerIngress, LoadBalancerStatus, Service, ServicePort, ServiceSpec,
ServiceStatus,
},
apimachinery::pkg::apis::meta::v1::{ObjectMeta, Time},
};
use kube::ResourceExt;
use std::{
collections::HashMap,
time::{Duration, Instant},
};

const RESYNC: Duration = Duration::from_secs(300);

fn service() -> Service {
Service {
metadata: ObjectMeta {
name: Some("web".to_string()),
uid: Some("uid-1".to_string()),
annotations: Some(
[(
consts::LB_LOCATION_LABEL_NAME.to_string(),
"fsn1".to_string(),
)]
.into(),
),
..Default::default()
},
spec: Some(ServiceSpec {
type_: Some("LoadBalancer".to_string()),
ports: Some(vec![ServicePort {
port: 80,
node_port: Some(30080),
..Default::default()
}]),
..Default::default()
}),
..Default::default()
}
}

fn fingerprint_of(svc: &Service, targets: &[&str], ports: &[(i32, i32)]) -> u64 {
let targets = targets.iter().map(ToString::to_string).collect::<Vec<_>>();
fingerprint(
svc,
&targets,
&ports.iter().copied().collect::<HashMap<_, _>>(),
)
}

fn print(svc: &Service) -> u64 {
fingerprint_of(svc, &["10.0.0.1"], &[(80, 30080)])
}

#[test]
fn an_unchanged_service_is_skipped_until_the_resync() {
let memo = ReconcileMemo::default();
let start = Instant::now();
memo.record("uid-1", 7, start, RESYNC);
assert_eq!(
memo.unchanged("uid-1", 7, start + Duration::from_secs(100)),
Some(Duration::from_secs(200))
);
assert_eq!(memo.unchanged("uid-1", 7, start + RESYNC), None);
assert_eq!(
memo.unchanged("uid-1", 7, start + RESYNC + Duration::from_secs(1)),
None
);
}

#[test]
fn a_shorter_validity_ends_the_skip_sooner() {
let memo = ReconcileMemo::default();
let start = Instant::now();
let retry = Duration::from_secs(30);
memo.record("uid-1", 7, start, retry);
assert_eq!(
memo.unchanged("uid-1", 7, start + Duration::from_secs(10)),
Some(Duration::from_secs(20))
);
assert_eq!(memo.unchanged("uid-1", 7, start + retry), None);
}

#[test]
fn a_changed_or_unknown_service_is_reconciled() {
let memo = ReconcileMemo::default();
let start = Instant::now();
assert_eq!(memo.unchanged("uid-1", 7, start), None);
memo.record("uid-1", 7, start, RESYNC);
assert_eq!(memo.unchanged("uid-1", 8, start), None);
assert_eq!(memo.unchanged("uid-2", 7, start), None);
}

#[test]
fn a_forgotten_service_is_reconciled() {
let memo = ReconcileMemo::default();
let start = Instant::now();
memo.record("uid-1", 7, start, RESYNC);
memo.forget("uid-1");
assert_eq!(memo.unchanged("uid-1", 7, start), None);
}

#[test]
fn a_new_record_restarts_the_resync() {
let memo = ReconcileMemo::default();
let start = Instant::now();
memo.record("uid-1", 7, start, RESYNC);
memo.record("uid-1", 8, start + RESYNC, RESYNC);
assert_eq!(memo.unchanged("uid-1", 8, start + RESYNC), Some(RESYNC));
}

fn settled(svc: &Service, succeeded: bool) -> bool {
let memo = ReconcileMemo::default();
let start = Instant::now();
memo.record("uid-1", 7, start, RESYNC);
memo.settle(svc, succeeded);
memo.unchanged("uid-1", 7, start).is_some()
}

#[test]
fn a_successful_reconcile_keeps_the_entry() {
assert!(settled(&service(), true));
}

#[test]
fn a_failed_reconcile_drops_the_entry() {
assert!(!settled(&service(), false));
}

#[test]
fn a_released_or_deleted_service_drops_the_entry() {
let mut svc = service();
svc.spec.as_mut().unwrap().type_ = Some("ClusterIP".to_string());
assert!(!settled(&svc, true));
let mut svc = service();
svc.spec.as_mut().unwrap().load_balancer_class = Some("other".to_string());
assert!(!settled(&svc, true));
let mut svc = service();
svc.metadata.deletion_timestamp = Some(Time(k8s_openapi::chrono::Utc::now()));
assert!(!settled(&svc, true));
}

#[test]
fn the_fingerprint_is_stable() {
assert_eq!(print(&service()), print(&service()));
}

#[test]
fn a_spec_change_changes_the_fingerprint() {
let mut svc = service();
svc.spec.as_mut().unwrap().external_traffic_policy = Some("Local".to_string());
assert_ne!(print(&svc), print(&service()));
}

#[test]
fn a_robotlb_annotation_changes_the_fingerprint() {
let mut svc = service();
svc.annotations_mut().insert(
consts::LB_LOCATION_LABEL_NAME.to_string(),
"nbg1".to_string(),
);
assert_ne!(print(&svc), print(&service()));
let mut svc = service();
svc.annotations_mut().insert(
consts::LB_PRIVATE_IP_LABEL_NAME.to_string(),
"10.0.0.9".to_string(),
);
assert_ne!(print(&svc), print(&service()));
}

#[test]
fn the_targets_change_the_fingerprint() {
let svc = service();
let one = fingerprint_of(&svc, &["10.0.0.1"], &[(80, 30080)]);
let two = fingerprint_of(&svc, &["10.0.0.1", "10.0.0.2"], &[(80, 30080)]);
let reordered = fingerprint_of(&svc, &["10.0.0.2", "10.0.0.1"], &[(80, 30080)]);
assert_ne!(one, two);
assert_eq!(two, reordered);
}

#[test]
fn the_ports_change_the_fingerprint() {
let svc = service();
let one = fingerprint_of(&svc, &["10.0.0.1"], &[(80, 30080)]);
let other = fingerprint_of(&svc, &["10.0.0.1"], &[(80, 30081)]);
assert_ne!(one, other);
}

#[test]
fn status_and_foreign_metadata_leave_the_fingerprint_alone() {
let mut svc = service();
svc.status = Some(ServiceStatus {
load_balancer: Some(LoadBalancerStatus {
ingress: Some(vec![LoadBalancerIngress {
ip: Some("192.0.2.1".to_string()),
..Default::default()
}]),
}),
..Default::default()
});
svc.metadata.resource_version = Some("42".to_string());
svc.annotations_mut().insert(
"metallb.io/ip-allocated-from-pool".to_string(),
"x".to_string(),
);
assert_eq!(print(&svc), print(&service()));
}
}
Loading