From 3821dc69ad47885b749b641ba965cdcc1bb20218 Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Wed, 30 Sep 2026 19:58:04 +0300 Subject: [PATCH] fix(controller): skip Hetzner calls when nothing robotlb acts on changed Every write to a Service starts a reconcile, including status writes by another LoadBalancer controller that also takes Services without a class. When that controller clears status.loadBalancer, robotlb writes it back, the other controller clears it again, and each round costs Hetzner requests until the project hits the rate limit. Services carry no metadata.generation, so a generation predicate cannot filter these events, and filtering the main stream needs an unstable kube-runtime feature. robotlb now remembers, per Service UID, a hash of the spec, the robotlb/ annotations and the computed targets and ports after each successful reconcile. A reconcile that finds the same hash before the next check is due makes no Hetzner requests and writes no status. Node and EndpointSlice changes still get through because they change the targets. An entry is valid until the requeue of the reconcile that recorded it, so while Hetzner refuses some targets events are skipped only until the 30-second retry. The entry is dropped after a failed reconcile, when the Service is released or deleted, and when it has no port to expose. The README now says which classes robotlb handles, that setting loadBalancerClass: robotlb keeps other controllers away, and that recreating a Service to add the class costs it its balancer and public IP. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- README.md | 4 +- src/main.rs | 29 ++++- src/memo.rs | 299 ++++++++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 327 insertions(+), 5 deletions(-) create mode 100644 src/memo.rs diff --git a/README.md b/README.md index 47f69c2..c321b75 100644 --- a/README.md +++ b/README.md @@ -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`: diff --git a/src/main.rs b/src/main.rs index aafa04e..b5cfe07 100644 --- a/src/main.rs +++ b/src/main.rs @@ -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; @@ -222,6 +223,7 @@ fn cluster_changes( #[derive(Clone)] pub struct CurrentContext { pub client: kube::Client, + pub memo: Arc, pub config: OperatorConfig, pub hcloud_config: HCloudConfig, pub rate_limit: Arc, @@ -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(), @@ -253,6 +256,7 @@ pub async fn reconcile_service( report_failure(&svc, &context, error).await; } } + context.memo.settle(&svc, result.is_ok()); result } @@ -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![]; @@ -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 diff --git a/src/memo.rs b/src/memo.rs new file mode 100644 index 0000000..806aa3d --- /dev/null +++ b/src/memo.rs @@ -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>, +} + +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 { + 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( + svc: &Service, + targets: &[String], + services: &HashMap, +) -> 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::>().hash(&mut hasher); + services + .iter() + .collect::>() + .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::>(); + fingerprint( + svc, + &targets, + &ports.iter().copied().collect::>(), + ) + } + + 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())); + } +}