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())); + } +}