From dc37dcfdb31dea7a47e40e5916fa1d24b24b63bf Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 29 Sep 2026 11:31:38 +0300 Subject: [PATCH 1/2] perf(controller): react to node and endpoint changes instead of polling Every LoadBalancer service was reconciled every 30 seconds, and each run cost one or two Hetzner API requests even when nothing changed. The short interval was there because only Services were watched, so a new node or a moved pod reached the balancer only through that requeue. Node changes that affect targets now reconcile every service. Heartbeat updates are ignored, since only labels, cordoning, readiness and addresses are compared. An endpoint slice change reconciles its service when it is a robotlb balancer with the Local traffic policy and the nodes its targets follow changed: all endpoint nodes when targets come from pods, ready endpoint nodes otherwise. Readiness flaps and address changes that leave those nodes alone cost nothing. The controller's watch mapper cannot tell a deletion from an update, so a second watch tracks which slices exist; a deleted slice with endpoints reconciles every service. The periodic resync is now ROBOTLB_RESYNC_INTERVAL, five minutes by default and at most a year. A reconcile that could not add some targets, for a reason other than the rate limit, is retried after 30 seconds like a failed one, so a temporarily rejected node does not wait for the resync. Pods that have finished or have no IP yet no longer count as targets. Neither is in the endpoint slice, so deleting them would leave their node a target until the resync. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- README.md | 4 + src/config.rs | 39 ++++ src/lb.rs | 11 +- src/main.rs | 454 +++++++++++++++++++++++++++++++++++++++++----- src/triggers.rs | 466 ++++++++++++++++++++++++++++++++++++++++++++++++ 5 files changed, 922 insertions(+), 52 deletions(-) create mode 100644 src/triggers.rs diff --git a/README.md b/README.md index 02c9aae..e348cba 100644 --- a/README.md +++ b/README.md @@ -35,6 +35,8 @@ 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 of their endpoints 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. + Target nodes are selected according to the service's `externalTrafficPolicy`: - `Cluster`, the Kubernetes default: every node of the cluster becomes a target, since kube-proxy forwards the traffic to a node that hosts a pod. Cordoned and not-ready nodes are left out, as they would only take up target slots. @@ -87,6 +89,8 @@ Options: Default load balancer proxy mode. If enabled, the load balancer will act as a proxy for the target servers. The default value is `false`. https://docs.hetzner.com/cloud/load-balancers/faq/#what-does-proxy-protocol-mean-and-should-i-enable-it [env: ROBOTLB_DEFAULT_LB_PROXY_MODE_ENABLED=] --ipv6-ingress Whether to enable IPv6 ingress for the load balancer. If enabled, the load balancer's IPv6 will be attached to the service as an external IP along with IPv4 [env: ROBOTLB_IPV6_INGRESS=] + --resync-interval + Seconds between reconciliations of a service that nothing changed. Node changes, and endpoint changes of Local services, trigger a reconciliation on their own; this interval bounds how long a change made to a balancer outside robotlb survives. A service whose balancer refused a target is retried within 30 seconds [env: ROBOTLB_RESYNC_INTERVAL=] [default: 300] --log-level [env: ROBOTLB_LOG_LEVEL=] [default: INFO] -h, --help diff --git a/src/config.rs b/src/config.rs index 7c3a25f..f2368ad 100644 --- a/src/config.rs +++ b/src/config.rs @@ -67,7 +67,46 @@ pub struct OperatorConfig { #[arg(long, env = "ROBOTLB_IPV6_INGRESS", default_value = "false")] pub ipv6_ingress: bool, + /// Seconds between reconciliations of a service that nothing changed. Node changes, + /// and endpoint changes of Local services, trigger a reconciliation on their own; + /// this interval bounds how long a change made to a balancer outside robotlb survives. + /// A service whose balancer refused a target is retried within 30 seconds. + #[arg( + long, + env = "ROBOTLB_RESYNC_INTERVAL", + default_value = "300", + value_parser = clap::value_parser!(u64).range(1..=31_536_000) + )] + pub resync_interval: u64, + // Log level of the operator. #[arg(long, env = "ROBOTLB_LOG_LEVEL", default_value = "INFO")] pub log_level: LevelFilter, } + +#[cfg(test)] +mod tests { + use super::OperatorConfig; + use clap::Parser; + + fn parse(resync: &str) -> Result { + OperatorConfig::try_parse_from([ + "robotlb", + "--hcloud-token", + "t", + "--resync-interval", + resync, + ]) + } + + // Zero would requeue every successful reconcile right away, spending Hetzner + // API requests on every service all the time. The controller's delay queue + // panics on delays past about two years. + #[test] + fn the_resync_interval_must_be_between_a_second_and_a_year() { + assert!(parse("0").is_err()); + assert!(parse("31536001").is_err()); + assert_eq!(parse("31536000").unwrap().resync_interval, 31_536_000); + assert_eq!(parse("1").unwrap().resync_interval, 1); + } +} diff --git a/src/lb.rs b/src/lb.rs index e1cc5f3..edd1478 100644 --- a/src/lb.rs +++ b/src/lb.rs @@ -172,14 +172,15 @@ impl LoadBalancer { /// Reconcile the load balancer to match the desired configuration. #[tracing::instrument(skip(self), fields(lb_name=self.name))] - pub async fn reconcile(&self) -> RobotLBResult { + /// Returns the balancer, and whether some of its targets could not be added. + pub async fn reconcile(&self) -> RobotLBResult<(hcloud::models::LoadBalancer, bool)> { let hcloud_balancer = self.get_or_create_hcloud_lb().await?; self.reconcile_algorithm(&hcloud_balancer).await?; self.reconcile_lb_type(&hcloud_balancer).await?; self.reconcile_network(&hcloud_balancer).await?; self.reconcile_services(&hcloud_balancer).await?; - self.reconcile_targets(&hcloud_balancer).await?; - Ok(hcloud_balancer) + let targets_missing = self.reconcile_targets(&hcloud_balancer).await?; + Ok((hcloud_balancer, targets_missing)) } /// Reconcile the services of the load balancer. @@ -299,7 +300,7 @@ impl LoadBalancer { async fn reconcile_targets( &self, hcloud_balancer: &hcloud::models::LoadBalancer, - ) -> RobotLBResult<()> { + ) -> RobotLBResult { let max_targets = usize::try_from(hcloud_balancer.load_balancer_type.max_targets).unwrap_or(usize::MAX); let planned = plan_targets(&self.targets, max_targets); @@ -386,7 +387,7 @@ impl LoadBalancer { last_error.unwrap_or_else(|| "no reason reported".to_string()), ))); } - Ok(()) + Ok(live < planned.len()) } /// Reconcile the load balancer algorithm. diff --git a/src/main.rs b/src/main.rs index edcfa73..9b53a3e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -33,7 +33,7 @@ use kube::{ runtime::{ controller::{self, Action}, events::{Event, EventType, Recorder, Reporter}, - watcher, Controller, + watcher, Controller, WatchStreamExt, }, Resource, ResourceExt, }; @@ -54,6 +54,7 @@ pub mod finalizers; pub mod label_filter; pub mod lb; pub mod rate_limit; +pub mod triggers; #[cfg(not(target_env = "msvc"))] #[global_allocator] @@ -73,7 +74,6 @@ async fn main() -> RobotLBResult<()> { tracing::info!("Starting robotlb operator v{}", env!("CARGO_PKG_VERSION")); let kube_client = kube::Client::try_default().await?; tracing::info!("Kube client is connected"); - watcher::Config::default(); let context = Arc::new(CurrentContext::new( kube_client.clone(), operator_config.clone(), @@ -81,41 +81,171 @@ async fn main() -> RobotLBResult<()> { )); tracing::info!("Starting the controller"); let token = operator_config.hcloud_token; - Controller::new( - kube::Api::::all(kube_client), + let dynamic_node_selector = operator_config.dynamic_node_selector; + let controller = Controller::new( + kube::Api::::all(kube_client.clone()), watcher::Config::default(), - ) - .run(reconcile_service, on_error, context) - .for_each(|reconcilation_result| { - let token = token.clone(); - async move { - match reconcilation_result { - Ok((service, _action)) => { - tracing::info!("Reconcilation of a service {} was successful", service.name); - } - // During reconcilation process, - // the controller has decided to skip the service. - Err(controller::Error::ReconcilerFailed(RobotLBError::SkipService, _)) => {} - Err(controller::Error::ReconcilerFailed( - error @ RobotLBError::RateLimited(_), - service, - )) => { - tracing::info!("Service {service}: {error}"); + ); + let services = controller.store(); + let slice_changes = Arc::new(triggers::SliceChanges::default()); + controller + .reconcile_all_on(cluster_changes(kube_client.clone(), slice_changes.clone())) + .watches( + kube::Api::::all(kube_client), + watcher::Config::default(), + move |slice| { + triggers::slice_service(&slice).filter(|service| { + let svc = services.get(service); + slice_triggers( + &slice, + svc.as_deref(), + dynamic_node_selector, + &slice_changes, + ) + }) + }, + ) + .run(reconcile_service, on_error, context) + .for_each(|reconcilation_result| { + let token = token.clone(); + async move { + match reconcilation_result { + Ok((service, _action)) => { + tracing::info!( + "Reconcilation of a service {} was successful", + service.name + ); + } + // During reconcilation process, + // the controller has decided to skip the service. + Err(controller::Error::ReconcilerFailed(RobotLBError::SkipService, _)) => {} + Err(controller::Error::ReconcilerFailed( + error @ RobotLBError::RateLimited(_), + service, + )) => { + tracing::info!("Service {service}: {error}"); + } + Err(controller::Error::ReconcilerFailed(error, service)) => { + tracing::error!( + "Reconcilation of service {service} failed: {}", + redact(&error.to_string(), &token) + ); + } + Err(error) => { + tracing::error!("Controller error: {}", redact(&error.to_string(), &token)); + } } - Err(controller::Error::ReconcilerFailed(error, service)) => { - tracing::error!( - "Reconcilation of service {service} failed: {}", - redact(&error.to_string(), &token) - ); + } + }) + .await; + Ok(()) +} + +/// The nodes of pods that may still serve traffic. A finished pod keeps its node +/// name, and neither it nor a pod without an IP is in the endpoint slice, so their +/// removal would otherwise go unnoticed until the resync. +fn pod_nodes(pods: &[Pod]) -> HashSet { + pods.iter() + .filter(|pod| { + let status = pod.status.as_ref(); + let phase = status.and_then(|status| status.phase.as_deref()); + let has_ip = status.and_then(|status| status.pod_ip.as_ref()).is_some(); + has_ip && !matches!(phase, Some("Succeeded" | "Failed")) + }) + .filter_map(|pod| pod.spec.as_ref()?.node_name.clone()) + .collect() +} + +/// The selector a service picks its pods with; without one its targets come from +/// its endpoint slices. +fn pod_selector(svc: &Service) -> Option> { + svc.spec + .as_ref() + .and_then(|spec| spec.selector.clone()) + .filter(|selector| !selector.is_empty()) +} + +fn is_robotlb_load_balancer(svc: &Service) -> bool { + let spec = svc.spec.as_ref(); + spec.and_then(|spec| spec.type_.as_deref()) == Some("LoadBalancer") + && spec + .and_then(|spec| spec.load_balancer_class.as_deref()) + .unwrap_or(consts::ROBOTLB_LB_CLASS) + == consts::ROBOTLB_LB_CLASS +} + +/// Whether a change of the slice should reconcile its service. Only services that +/// target the nodes of their endpoints care where pods run. +fn slice_triggers( + slice: &EndpointSlice, + svc: Option<&Service>, + dynamic_node_selector: bool, + changes: &triggers::SliceChanges, +) -> bool { + let observed = svc.is_some_and(|svc| { + is_robotlb_load_balancer(svc) + && node_source(svc, dynamic_node_selector) == NodeSource::ServiceEndpoints + }); + if !observed { + changes.forget(slice); + return false; + } + let follows = if svc.and_then(pod_selector).is_some() { + triggers::Follows::AllNodes + } else { + triggers::Follows::ReadyNodes + }; + changes.observe(slice, follows) +} + +/// A signal each time nodes change in a way that can change balancer targets, or an +/// endpoint slice with endpoints is deleted. +fn cluster_changes( + client: kube::Client, + slices: Arc, +) -> impl futures::Stream + Send + Sync { + let (sender, receiver) = futures::channel::mpsc::unbounded(); + let node_sender = sender.clone(); + let node_client = client.clone(); + tokio::spawn(async move { + let mut changes = triggers::NodeChanges::default(); + let mut events = std::pin::pin!(watcher::watcher( + kube::Api::::all(node_client), + watcher::Config::default() + ) + .default_backoff()); + while let Some(event) = events.next().await { + match event { + Ok(event) => { + if changes.observe(&event) && node_sender.unbounded_send(()).is_err() { + return; + } } - Err(error) => { - tracing::error!("Controller error: {}", redact(&error.to_string(), &token)); + Err(error) => tracing::warn!("Node watch failed: {error}"), + } + } + }); + // A deleted slice triggers every service, since only the unstable + // `Controller::reconcile_on` can target one. Once kube-runtime stabilizes it, + // send the slice's service there instead and drop this watch's reconcile-all. + tokio::spawn(async move { + let mut events = std::pin::pin!(watcher::metadata_watcher( + kube::Api::::all(client), + watcher::Config::default() + ) + .default_backoff()); + while let Some(event) = events.next().await { + match event { + Ok(event) => { + if slices.track(&event) && sender.unbounded_send(()).is_err() { + return; + } } + Err(error) => tracing::warn!("Endpoint slice watch failed: {error}"), } } - }) - .await; - Ok(()) + }); + receiver } #[derive(Clone)] @@ -264,12 +394,7 @@ async fn get_nodes_dynamically( .unwrap_or_else(|| context.client.default_namespace()), ); - let Some(pod_selector) = svc - .spec - .as_ref() - .and_then(|spec| spec.selector.clone()) - .filter(|s| !s.is_empty()) - else { + let Some(pod_selector) = pod_selector(svc) else { tracing::info!( "Service has no selector, falling back to EndpointSlice-based node discovery" ); @@ -289,11 +414,7 @@ async fn get_nodes_dynamically( }) .await?; - let target_nodes = pods - .iter() - .map(|pod| pod.spec.clone().unwrap_or_default().node_name) - .flatten() - .collect::>(); + let target_nodes = pod_nodes(&pods.items); let nodes_api = kube::Api::::all(context.client.clone()); let nodes = nodes_api @@ -528,10 +649,12 @@ 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?; - return Ok(Action::requeue(Duration::from_secs(30))); + return Ok(Action::requeue(Duration::from_secs( + context.config.resync_interval, + ))); } - let hcloud_lb = lb.reconcile().await?; + let (hcloud_lb, targets_missing) = lb.reconcile().await?; let mut ingress = vec![]; @@ -572,7 +695,20 @@ pub async fn reconcile_load_balancer( .await?; } - Ok(Action::requeue(Duration::from_secs(30))) + Ok(Action::requeue(success_requeue( + targets_missing, + context.config.resync_interval, + ))) +} + +/// A target Hetzner rejected, for a reason other than the rate limit, may be accepted +/// on the next try, so it is retried as soon as a failed reconcile would be. +const fn success_requeue(targets_missing: bool, resync_interval: u64) -> Duration { + if targets_missing && resync_interval > 30 { + Duration::from_secs(30) + } else { + Duration::from_secs(resync_interval) + } } /// Drop the external IP a service advertises, so that nothing keeps sending @@ -631,7 +767,8 @@ fn error_action( mod tests { use super::{ collect_lb_services, consts, error_action, event_note, is_excluded_from_lb, - is_lb_eligible_node, is_local_traffic_policy, node_source, publishes_event, NodeSource, + is_lb_eligible_node, is_local_traffic_policy, node_source, pod_nodes, publishes_event, + slice_triggers, success_requeue, NodeSource, }; use k8s_openapi::{ api::core::v1::{ @@ -639,7 +776,7 @@ mod tests { }, apimachinery::pkg::apis::meta::v1::ObjectMeta, }; - use std::collections::BTreeMap; + use std::collections::{BTreeMap, HashSet}; fn service(spec: ServiceSpec) -> Service { Service { @@ -908,4 +1045,227 @@ mod tests { "boom".to_string() ))); } + + fn endpoint_slice(node: &str) -> k8s_openapi::api::discovery::v1::EndpointSlice { + k8s_openapi::api::discovery::v1::EndpointSlice { + metadata: ObjectMeta { + name: Some("web-abcde".to_string()), + namespace: Some("shop".to_string()), + ..Default::default() + }, + endpoints: vec![k8s_openapi::api::discovery::v1::Endpoint { + node_name: Some(node.to_string()), + ..Default::default() + }], + ..Default::default() + } + } + + fn live( + slice: &k8s_openapi::api::discovery::v1::EndpointSlice, + ) -> crate::triggers::SliceChanges { + use kube::runtime::watcher::Event; + let changes = crate::triggers::SliceChanges::default(); + changes.track(&Event::::Init); + changes.track(&Event::InitApply(slice.clone())); + changes.track(&Event::::InitDone); + changes + } + + fn with_policy(policy: &str) -> Service { + service(ServiceSpec { + type_: Some("LoadBalancer".to_string()), + external_traffic_policy: Some(policy.to_string()), + ..Default::default() + }) + } + + #[test] + fn a_slice_of_a_local_service_triggers_only_when_its_nodes_change() { + let slice = endpoint_slice("a"); + let changes = live(&slice); + let svc = with_policy("Local"); + assert!(slice_triggers(&slice, Some(&svc), true, &changes)); + assert!(!slice_triggers(&slice, Some(&svc), true, &changes)); + assert!(slice_triggers( + &endpoint_slice("b"), + Some(&svc), + true, + &changes + )); + } + + #[test] + fn a_slice_of_a_service_that_ignores_endpoints_triggers_nothing_and_is_forgotten() { + let slice = endpoint_slice("a"); + let changes = live(&slice); + let local = with_policy("Local"); + slice_triggers(&slice, Some(&local), true, &changes); + assert!(!slice_triggers( + &slice, + Some(&with_policy("Cluster")), + true, + &changes + )); + assert!(!slice_triggers(&slice, Some(&local), false, &changes)); + assert!(slice_triggers(&slice, Some(&local), true, &changes)); + } + + #[test] + fn a_slice_without_a_known_service_triggers_nothing() { + let slice = endpoint_slice("a"); + assert!(!slice_triggers(&slice, None, true, &live(&slice))); + } + + fn pod(node: &str, phase: &str) -> k8s_openapi::api::core::v1::Pod { + k8s_openapi::api::core::v1::Pod { + spec: Some(k8s_openapi::api::core::v1::PodSpec { + node_name: Some(node.to_string()), + ..Default::default() + }), + status: Some(k8s_openapi::api::core::v1::PodStatus { + phase: Some(phase.to_string()), + pod_ip: Some("10.0.0.1".to_string()), + ..Default::default() + }), + ..Default::default() + } + } + + // An evicted pod keeps its node name, but nothing on that node serves the + // service any more, and its later deletion does not touch the endpoint slice. + #[test] + fn pods_that_finished_are_not_targets() { + let pods = [ + pod("a", "Running"), + pod("b", "Failed"), + pod("c", "Succeeded"), + pod("d", "Pending"), + ]; + assert_eq!( + pod_nodes(&pods), + HashSet::from(["a".to_string(), "d".to_string()]) + ); + } + + // A pod without an IP is not in the endpoint slice either, so its deletion + // would go unnoticed until the resync. + #[test] + fn pods_without_an_ip_are_not_targets() { + let mut scheduled = pod("a", "Pending"); + scheduled.status.as_mut().unwrap().pod_ip = None; + assert!(pod_nodes(&[scheduled]).is_empty()); + } + + fn flapping_slice(ready: bool) -> k8s_openapi::api::discovery::v1::EndpointSlice { + let mut slice = endpoint_slice("a"); + slice.endpoints[0].conditions = Some(k8s_openapi::api::discovery::v1::EndpointConditions { + ready: Some(ready), + ..Default::default() + }); + slice + } + + // Targets of a service with a selector come from its pods, so readiness of its + // endpoints moves nothing; without a selector they come from ready endpoints. + #[test] + fn a_readiness_flap_triggers_only_services_whose_targets_follow_readiness() { + let with_selector = service(ServiceSpec { + type_: Some("LoadBalancer".to_string()), + external_traffic_policy: Some("Local".to_string()), + selector: Some(BTreeMap::from([("app".to_string(), "web".to_string())])), + ..Default::default() + }); + let ready = flapping_slice(true); + let changes = live(&ready); + slice_triggers(&ready, Some(&with_selector), true, &changes); + assert!(!slice_triggers( + &flapping_slice(false), + Some(&with_selector), + true, + &changes + )); + + let without_selector = with_policy("Local"); + let changes = live(&ready); + slice_triggers(&ready, Some(&without_selector), true, &changes); + assert!(slice_triggers( + &flapping_slice(false), + Some(&without_selector), + true, + &changes + )); + } + + // Deleting a slice with recorded nodes reconciles every balancer, which a service + // without a robotlb balancer must not cost. + #[test] + fn slices_of_services_without_a_robotlb_balancer_are_not_observed() { + let slice = endpoint_slice("a"); + for spec in [ + ServiceSpec { + type_: Some("NodePort".to_string()), + external_traffic_policy: Some("Local".to_string()), + ..Default::default() + }, + ServiceSpec { + type_: Some("LoadBalancer".to_string()), + load_balancer_class: Some("example.com/other".to_string()), + external_traffic_policy: Some("Local".to_string()), + ..Default::default() + }, + ] { + let changes = live(&slice); + assert!(!slice_triggers( + &slice, + Some(&service(spec)), + true, + &changes + )); + assert!(!changes.track(&kube::runtime::watcher::Event::Delete(slice.clone()))); + } + } + + // A service that gains or loses its selector switches which nodes its targets + // follow, and the nodes recorded under the old rule must not hide a change. + #[test] + fn a_selector_change_does_not_hide_the_next_endpoint_change() { + let mut both = endpoint_slice("a"); + both.endpoints + .push(k8s_openapi::api::discovery::v1::Endpoint { + node_name: Some("b".to_string()), + conditions: Some(k8s_openapi::api::discovery::v1::EndpointConditions { + ready: Some(false), + ..Default::default() + }), + ..Default::default() + }); + let with_selector = service(ServiceSpec { + type_: Some("LoadBalancer".to_string()), + external_traffic_policy: Some("Local".to_string()), + selector: Some(BTreeMap::from([("app".to_string(), "web".to_string())])), + ..Default::default() + }); + let changes = live(&both); + slice_triggers(&both, Some(&with_selector), true, &changes); + // Under the new rule the nodes happen to equal what the old rule recorded. + let mut both_ready = both.clone(); + both_ready.endpoints[1].conditions = None; + assert!(slice_triggers( + &both_ready, + Some(&with_policy("Local")), + true, + &changes + )); + } + + // A node whose target could not be added would otherwise wait for the resync, + // since a partly successful reconcile still counts as a success. + #[test] + fn missing_targets_are_retried_like_an_error() { + use std::time::Duration; + assert_eq!(success_requeue(true, 300), Duration::from_secs(30)); + assert_eq!(success_requeue(false, 300), Duration::from_secs(300)); + assert_eq!(success_requeue(true, 10), Duration::from_secs(10)); + } } diff --git a/src/triggers.rs b/src/triggers.rs new file mode 100644 index 0000000..31d24c1 --- /dev/null +++ b/src/triggers.rs @@ -0,0 +1,466 @@ +use k8s_openapi::api::{ + core::v1::{Node, Service}, + discovery::v1::EndpointSlice, +}; +use kube::{ + runtime::{reflector::ObjectRef, watcher::Event}, + ResourceExt, +}; +use std::{ + collections::{hash_map::DefaultHasher, BTreeSet, HashMap, HashSet}, + hash::{Hash, Hasher}, + sync::Mutex, +}; + +/// Tells which node events change what a balancer may target. +/// +/// Kubelets rewrite their node's status on every heartbeat, and reconciling every +/// service on each of those would call Hetzner for every service every few minutes +/// per node. +#[derive(Debug, Default)] +pub struct NodeChanges { + seen: HashMap, + relisted: HashMap, + synced: bool, +} + +impl NodeChanges { + /// Whether the event calls for reconciling every service. + pub fn observe(&mut self, event: &Event) -> bool { + match event { + Event::Init => { + self.relisted.clear(); + false + } + Event::InitApply(node) => { + self.relisted.insert(node.name_any(), fingerprint(node)); + false + } + Event::InitDone => { + // The first listing counts as a change too: a node may have changed + // after the controller's first reconciles listed it. + let changed = !self.synced || self.relisted != self.seen; + self.seen = std::mem::take(&mut self.relisted); + self.synced = true; + changed + } + Event::Apply(node) => { + let print = fingerprint(node); + self.seen.insert(node.name_any(), print) != Some(print) + } + Event::Delete(node) => self.seen.remove(&node.name_any()).is_some(), + } + } +} + +/// The node fields target selection reads: labels, cordoning, readiness and addresses. +fn fingerprint(node: &Node) -> u64 { + let mut hasher = DefaultHasher::new(); + node.labels().hash(&mut hasher); + node.spec + .as_ref() + .and_then(|spec| spec.unschedulable) + .hash(&mut hasher); + if let Some(status) = &node.status { + status + .conditions + .iter() + .flatten() + .find(|condition| condition.type_ == "Ready") + .map(|condition| condition.status.as_str()) + .hash(&mut hasher); + for address in status.addresses.iter().flatten() { + (address.type_.as_str(), address.address.as_str()).hash(&mut hasher); + } + } + hasher.finish() +} + +type SliceKey = (String, String); + +/// Which endpoint nodes a service's targets follow. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Follows { + /// Targets come from the service's pods, which ignore readiness. + AllNodes, + /// Targets come from the ready endpoints. + ReadyNodes, +} + +/// Tells which endpoint slice changes move balancer targets. +/// +/// A slice is rewritten on every readiness flip and address change, but only the +/// nodes of its endpoints matter to a balancer. Reacting to every rewrite would call +/// Hetzner for a busy service every few seconds. +/// +/// The controller's watch mapper cannot tell a deletion from an update, so a +/// separate watch tracks which slices exist and reports deletions. +#[derive(Debug, Default)] +pub struct SliceChanges { + state: Mutex, +} + +#[derive(Debug, Default)] +struct SliceState { + live: HashSet, + relisted: HashSet, + nodes: HashMap)>, +} + +impl SliceChanges { + /// Whether the nodes of the slice changed. A slice the watch does not know to + /// exist, just created or already deleted, always counts as changed. + pub fn observe(&self, slice: &EndpointSlice, follows: Follows) -> bool { + let key = slice_key(slice); + let nodes = slice_nodes(slice, follows); + let mut state = self.lock(); + if !state.live.contains(&key) { + return true; + } + let recorded = (follows, nodes); + state.nodes.insert(key, recorded.clone()) != Some(recorded) + } + + /// Drop what is known about the slice's nodes, for a service no longer observed. + pub fn forget(&self, slice: &EndpointSlice) { + self.lock().nodes.remove(&slice_key(slice)); + } + + /// Record a watch event. Returns whether a slice with endpoints was deleted, + /// which calls for reconciling every service. + /// Generic over the object, so the watch can carry only slice metadata. + pub fn track(&self, event: &Event) -> bool { + let mut state = self.lock(); + match event { + Event::Init => { + state.relisted.clear(); + false + } + Event::InitApply(slice) => { + state.relisted.insert(slice_key(slice)); + false + } + Event::InitDone => { + let live = std::mem::take(&mut state.relisted); + let mut dropped_endpoints = false; + state.nodes.retain(|key, (_, nodes)| { + let kept = live.contains(key); + dropped_endpoints |= !kept && !nodes.is_empty(); + kept + }); + state.live = live; + dropped_endpoints + } + Event::Apply(slice) => { + state.live.insert(slice_key(slice)); + false + } + Event::Delete(slice) => { + let key = slice_key(slice); + state.live.remove(&key); + state + .nodes + .remove(&key) + .is_some_and(|(_, nodes)| !nodes.is_empty()) + } + } + } + + fn lock(&self) -> std::sync::MutexGuard<'_, SliceState> { + self.state + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } +} + +fn slice_key(slice: &K) -> SliceKey { + (slice.namespace().unwrap_or_default(), slice.name_any()) +} + +fn slice_nodes(slice: &EndpointSlice, follows: Follows) -> BTreeSet { + slice + .endpoints + .iter() + // The API treats a missing ready condition as ready. + .filter(|endpoint| { + follows == Follows::AllNodes + || endpoint.conditions.as_ref().and_then(|c| c.ready) != Some(false) + }) + .filter_map(|endpoint| endpoint.node_name.clone()) + .collect() +} + +#[must_use] +pub fn slice_service(slice: &EndpointSlice) -> Option> { + let name = slice.labels().get("kubernetes.io/service-name")?; + let namespace = slice.namespace()?; + Some(ObjectRef::new(name).within(&namespace)) +} + +#[cfg(test)] +mod tests { + use super::{slice_service, Follows, NodeChanges, SliceChanges}; + use k8s_openapi::{ + api::{ + core::v1::{Node, NodeAddress, NodeCondition, NodeSpec, NodeStatus}, + discovery::v1::{Endpoint, EndpointConditions, EndpointSlice}, + }, + apimachinery::pkg::apis::meta::v1::ObjectMeta, + }; + use kube::runtime::{reflector::ObjectRef, watcher::Event}; + use std::collections::BTreeMap; + + fn node(name: &str, ready: &str) -> Node { + Node { + metadata: ObjectMeta { + name: Some(name.to_string()), + ..Default::default() + }, + spec: Some(NodeSpec::default()), + status: Some(NodeStatus { + addresses: Some(vec![NodeAddress { + type_: "InternalIP".to_string(), + address: "192.0.2.10".to_string(), + }]), + conditions: Some(vec![NodeCondition { + type_: "Ready".to_string(), + status: ready.to_string(), + ..Default::default() + }]), + ..Default::default() + }), + } + } + + fn synced(nodes: &[Node]) -> NodeChanges { + let mut changes = NodeChanges::default(); + assert!(!changes.observe(&Event::Init)); + for node in nodes { + assert!(!changes.observe(&Event::InitApply(node.clone()))); + } + changes.observe(&Event::InitDone); + changes + } + + // A node that changed between the controller's first reconciles and this watch's + // first listing would otherwise only be picked up by the resync. + #[test] + fn the_first_listing_triggers_a_reconcile() { + let mut changes = NodeChanges::default(); + assert!(!changes.observe(&Event::Init)); + assert!(!changes.observe(&Event::InitApply(node("a", "True")))); + assert!(changes.observe(&Event::InitDone)); + } + + #[test] + fn a_heartbeat_triggers_nothing() { + let mut changes = synced(&[node("a", "True")]); + let mut beat = node("a", "True"); + if let Some(conditions) = beat.status.as_mut().and_then(|s| s.conditions.as_mut()) { + conditions[0].message = Some("kubelet is posting ready status".to_string()); + } + assert!(!changes.observe(&Event::Apply(beat))); + } + + #[test] + fn readiness_labels_addresses_and_cordoning_trigger_a_reconcile() { + let mut changes = synced(&[node("a", "True")]); + assert!(changes.observe(&Event::Apply(node("a", "False")))); + assert!(!changes.observe(&Event::Apply(node("a", "False")))); + + let mut labelled = node("a", "False"); + labelled.metadata.labels = Some(BTreeMap::from([( + "node.kubernetes.io/exclude-from-external-load-balancers".to_string(), + String::new(), + )])); + assert!(changes.observe(&Event::Apply(labelled.clone()))); + + let mut moved = labelled; + moved.status.as_mut().unwrap().addresses.as_mut().unwrap()[0].address = + "192.0.2.11".to_string(); + assert!(changes.observe(&Event::Apply(moved.clone()))); + + let mut cordoned = moved; + cordoned.spec.as_mut().unwrap().unschedulable = Some(true); + assert!(changes.observe(&Event::Apply(cordoned))); + } + + #[test] + fn added_and_removed_nodes_trigger_a_reconcile() { + let mut changes = synced(&[node("a", "True")]); + assert!(changes.observe(&Event::Apply(node("b", "True")))); + assert!(changes.observe(&Event::Delete(node("b", "True")))); + assert!(!changes.observe(&Event::Delete(node("unknown", "True")))); + } + + #[test] + fn a_node_gone_from_a_relisting_triggers_a_reconcile() { + let mut changes = synced(&[node("a", "True"), node("b", "True")]); + assert!(!changes.observe(&Event::Init)); + assert!(!changes.observe(&Event::InitApply(node("a", "True")))); + assert!(changes.observe(&Event::InitDone)); + } + + #[test] + fn a_relisting_triggers_only_when_nodes_changed_meanwhile() { + let mut changes = synced(&[node("a", "True")]); + assert!(!changes.observe(&Event::Init)); + assert!(!changes.observe(&Event::InitApply(node("a", "True")))); + assert!(!changes.observe(&Event::InitDone)); + + assert!(!changes.observe(&Event::Init)); + assert!(!changes.observe(&Event::InitApply(node("a", "False")))); + assert!(changes.observe(&Event::InitDone)); + } + + #[test] + fn an_endpoint_slice_points_at_its_service() { + let slice = EndpointSlice { + metadata: ObjectMeta { + name: Some("web-abcde".to_string()), + namespace: Some("shop".to_string()), + labels: Some(BTreeMap::from([( + "kubernetes.io/service-name".to_string(), + "web".to_string(), + )])), + ..Default::default() + }, + ..Default::default() + }; + assert_eq!( + slice_service(&slice), + Some(ObjectRef::new("web").within("shop")) + ); + } + + #[test] + fn a_slice_without_a_service_points_nowhere() { + assert_eq!(slice_service(&EndpointSlice::default()), None); + } + + fn slice(endpoints: &[(&str, &str, bool)]) -> EndpointSlice { + EndpointSlice { + metadata: ObjectMeta { + name: Some("web-abcde".to_string()), + namespace: Some("shop".to_string()), + ..Default::default() + }, + endpoints: endpoints + .iter() + .map(|(address, node, ready)| Endpoint { + addresses: vec![(*address).to_string()], + node_name: Some((*node).to_string()), + conditions: Some(EndpointConditions { + ready: Some(*ready), + ..Default::default() + }), + ..Default::default() + }) + .collect(), + ..Default::default() + } + } + + fn tracked(initial: &EndpointSlice, follows: Follows) -> SliceChanges { + let changes = SliceChanges::default(); + assert!(!changes.track(&Event::::Init)); + assert!(!changes.track(&Event::InitApply(initial.clone()))); + assert!(!changes.track(&Event::::InitDone)); + changes.observe(initial, follows); + changes + } + + #[test] + fn a_slice_created_after_the_listing_is_tracked_like_a_listed_one() { + let changes = tracked(&slice(&[]), Follows::ReadyNodes); + let created = slice(&[("10.0.0.1", "a", true)]); + let mut created = created; + created.metadata.name = Some("web-fghij".to_string()); + assert!(!changes.track(&Event::Apply(created.clone()))); + assert!(changes.observe(&created, Follows::ReadyNodes)); + assert!(!changes.observe(&created, Follows::ReadyNodes)); + assert!(changes.track(&Event::Delete(created))); + } + + #[test] + fn a_slice_not_known_to_exist_triggers_a_reconcile() { + let changes = SliceChanges::default(); + assert!(changes.observe(&slice(&[("10.0.0.1", "a", true)]), Follows::ReadyNodes)); + } + + // Replicas flapping readiness or changing addresses on the same nodes rewrite the + // slice all the time, and none of that moves a balancer target. + #[test] + fn endpoint_churn_on_the_same_nodes_triggers_nothing() { + for follows in [Follows::AllNodes, Follows::ReadyNodes] { + let changes = tracked( + &slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "a", true)]), + follows, + ); + let moved = slice(&[("10.0.0.3", "a", true), ("10.0.0.2", "a", true)]); + assert!(!changes.observe(&moved, follows)); + let flapped = slice(&[("10.0.0.3", "a", false), ("10.0.0.2", "a", true)]); + assert!(!changes.observe(&flapped, follows)); + } + } + + #[test] + fn endpoint_targets_follow_the_last_ready_endpoint_of_a_node() { + let changes = tracked(&slice(&[("10.0.0.1", "a", true)]), Follows::ReadyNodes); + let added = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", true)]); + assert!(changes.observe(&added, Follows::ReadyNodes)); + let unready = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", false)]); + assert!(changes.observe(&unready, Follows::ReadyNodes)); + } + + // Pod-based targets ignore readiness, so a flap costs nothing there, while an + // endpoint leaving the slice, terminating or not, moves a target. + #[test] + fn pod_targets_ignore_readiness_and_follow_endpoints_leaving() { + let both = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", true)]); + let changes = tracked(&both, Follows::AllNodes); + let unready = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", false)]); + assert!(!changes.observe(&unready, Follows::AllNodes)); + assert!(changes.observe(&slice(&[("10.0.0.1", "a", true)]), Follows::AllNodes)); + } + + #[test] + fn deleting_a_slice_with_endpoints_triggers_a_reconcile_of_everything() { + let initial = slice(&[("10.0.0.1", "a", true)]); + let changes = tracked(&initial, Follows::ReadyNodes); + assert!(changes.track(&Event::Delete(initial.clone()))); + // The mapper may see the deletion after the watch did; the slice is gone by + // then, so it reconciles the service rather than recording the slice again. + assert!(changes.observe(&initial, Follows::ReadyNodes)); + assert!(!changes.track(&Event::Delete(initial))); + } + + // A slice deleted while the watch was reconnecting is only noticed by its + // absence from the new listing, and counts as a deletion. + #[test] + fn a_relisting_forgets_slices_that_are_gone() { + let initial = slice(&[("10.0.0.1", "a", true)]); + let changes = tracked(&initial, Follows::ReadyNodes); + assert!(!changes.track(&Event::::Init)); + assert!(changes.track(&Event::::InitDone)); + assert!(changes.observe(&initial, Follows::ReadyNodes)); + } + + // A service that stops using the Local policy stops being observed, and its old + // record must not hide a change once it uses the policy again. + #[test] + fn a_forgotten_slice_counts_as_changed() { + let initial = slice(&[("10.0.0.1", "a", true)]); + let changes = tracked(&initial, Follows::ReadyNodes); + changes.forget(&initial); + assert!(changes.observe(&initial, Follows::ReadyNodes)); + } + + #[test] + fn a_relisting_that_drops_only_empty_slices_triggers_nothing() { + let empty = slice(&[]); + let changes = tracked(&empty, Follows::ReadyNodes); + assert!(!changes.track(&Event::::Init)); + assert!(!changes.track(&Event::::InitDone)); + } +} From f245740bb2d9ce13d5be4df87cede5e90f9e71aa Mon Sep 17 00:00:00 2001 From: Aleksei Sviridkin Date: Tue, 6 Oct 2026 11:12:30 +0300 Subject: [PATCH 2/2] fix(controller): take Local targets from the nodes kube-proxy serves (#56) A Local service with a selector took its targets from the pods the selector matched. That kept nodes where kube-proxy had no endpoint to send traffic to, such as the node of a pod that is not ready yet, and so nodes that fail the balancer's health check. A cordoned or not-ready node running a matching pod stayed a target too, unlike under the Cluster policy, with nothing saying whether that was intended. Every Local service now takes its targets from its EndpointSlices, the path already used for services without a selector. A node is a target when kube-proxy serves Local traffic from it: it has a ready endpoint, or a terminating endpoint that still serves, which kube-proxy falls back to while the node has no ready one. Missing conditions are read the way kube-proxy reads them. Targets and the slice trigger share this filter, so the trigger no longer needs a separate pod-based mode. robotlb no longer reads pods at all, and the chart's default role drops its access to them. Assisted-by: LLM Signed-off-by: Aleksei Sviridkin --- README.md | 11 +-- helm/values.yaml | 2 +- src/config.rs | 2 +- src/main.rs | 196 ++++------------------------------------------- src/triggers.rs | 160 +++++++++++++++++++++++++------------- 5 files changed, 129 insertions(+), 242 deletions(-) diff --git a/README.md b/README.md index e348cba..47f69c2 100644 --- a/README.md +++ b/README.md @@ -35,12 +35,12 @@ 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 of their endpoints 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. +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. Target nodes are selected according to the service's `externalTrafficPolicy`: - `Cluster`, the Kubernetes default: every node of the cluster becomes a target, since kube-proxy forwards the traffic to a node that hosts a pod. Cordoned and not-ready nodes are left out, as they would only take up target slots. -- `Local`: only the nodes where the service's target pods run, found through the service selector, or through the service's `EndpointSlice` resources when it has no selector. +- `Local`: only the nodes kube-proxy serves the traffic from, read from the service's `EndpointSlice` resources whether or not the service has a selector. These are the nodes with a ready endpoint, or with a terminating endpoint that still serves while the node has no ready one. Endpoint conditions decide, not whether a pod runs on the node: a node whose pods are not ready is left out, and a cordoned or not-ready node stays a target while it has such an endpoint. Nodes labelled `node.kubernetes.io/exclude-from-external-load-balancers`, which kubeadm puts on control-plane nodes, stay out of every balancer under either policy. @@ -72,7 +72,7 @@ Options: --default-network Default network to use for load balancers. If not set, then only network from the service annotation will be used [env: ROBOTLB_DEFAULT_NETWORK=] --dynamic-node-selector - If enabled, the operator will try to find target nodes based on where the target pods are actually deployed. If disabled, the operator will try to find target nodes based on the node selector [env: ROBOTLB_DYNAMIC_NODE_SELECTOR=] + If enabled, the operator will try to find target nodes based on the service's traffic policy, and under `Local` on the nodes serving the service's endpoints. If disabled, the operator will try to find target nodes based on the node selector [env: ROBOTLB_DYNAMIC_NODE_SELECTOR=] --default-lb-retries Default load balancer healthcheck retries cound [env: ROBOTLB_DEFAULT_LB_RETRIES=] [default: 3] --default-lb-timeout @@ -146,8 +146,9 @@ metadata: robotlb/balancer-type: "lb11" spec: type: LoadBalancer - # If dynamic node selector is enabled, nodes will be found - # using this property. + # The selector fills the service's endpoints; with + # externalTrafficPolicy: Local and the dynamic node selector + # enabled, the nodes serving them become the targets. selector: app: target ports: diff --git a/helm/values.yaml b/helm/values.yaml index 865062a..e696f57 100644 --- a/helm/values.yaml +++ b/helm/values.yaml @@ -34,7 +34,7 @@ serviceAccount: resources: [services, services/status] verbs: [get, list, patch, update, watch] - apiGroups: [""] - resources: [nodes, pods] + resources: [nodes] verbs: [get, list, watch] - apiGroups: [discovery.k8s.io] resources: [endpointslices] diff --git a/src/config.rs b/src/config.rs index f2368ad..6ce213b 100644 --- a/src/config.rs +++ b/src/config.rs @@ -12,7 +12,7 @@ pub struct OperatorConfig { #[arg(long, env = "ROBOTLB_DEFAULT_NETWORK", default_value = None)] pub default_network: Option, - /// If enabled, the operator will try to find target nodes based on where the target pods are actually deployed. + /// If enabled, the operator will try to find target nodes based on the service's traffic policy, and under `Local` on the nodes serving the service's endpoints. /// If disabled, the operator will try to find target nodes based on the node selector. #[arg(long, env = "ROBOTLB_DYNAMIC_NODE_SELECTOR", default_value = "true")] pub dynamic_node_selector: bool, diff --git a/src/main.rs b/src/main.rs index 9b53a3e..aafa04e 100644 --- a/src/main.rs +++ b/src/main.rs @@ -23,7 +23,7 @@ use futures::StreamExt; use hcloud::apis::configuration::Configuration as HCloudConfig; use k8s_openapi::{ api::{ - core::v1::{Node, Pod, Service}, + core::v1::{Node, Service}, discovery::v1::EndpointSlice, }, serde_json::json, @@ -141,30 +141,6 @@ async fn main() -> RobotLBResult<()> { Ok(()) } -/// The nodes of pods that may still serve traffic. A finished pod keeps its node -/// name, and neither it nor a pod without an IP is in the endpoint slice, so their -/// removal would otherwise go unnoticed until the resync. -fn pod_nodes(pods: &[Pod]) -> HashSet { - pods.iter() - .filter(|pod| { - let status = pod.status.as_ref(); - let phase = status.and_then(|status| status.phase.as_deref()); - let has_ip = status.and_then(|status| status.pod_ip.as_ref()).is_some(); - has_ip && !matches!(phase, Some("Succeeded" | "Failed")) - }) - .filter_map(|pod| pod.spec.as_ref()?.node_name.clone()) - .collect() -} - -/// The selector a service picks its pods with; without one its targets come from -/// its endpoint slices. -fn pod_selector(svc: &Service) -> Option> { - svc.spec - .as_ref() - .and_then(|spec| spec.selector.clone()) - .filter(|selector| !selector.is_empty()) -} - fn is_robotlb_load_balancer(svc: &Service) -> bool { let spec = svc.spec.as_ref(); spec.and_then(|spec| spec.type_.as_deref()) == Some("LoadBalancer") @@ -175,7 +151,7 @@ fn is_robotlb_load_balancer(svc: &Service) -> bool { } /// Whether a change of the slice should reconcile its service. Only services that -/// target the nodes of their endpoints care where pods run. +/// target the nodes of their endpoints care where those run. fn slice_triggers( slice: &EndpointSlice, svc: Option<&Service>, @@ -190,12 +166,7 @@ fn slice_triggers( changes.forget(slice); return false; } - let follows = if svc.and_then(pod_selector).is_some() { - triggers::Follows::AllNodes - } else { - triggers::Follows::ReadyNodes - }; - changes.observe(slice, follows) + changes.observe(slice) } /// A signal each time nodes change in a way that can change balancer targets, or an @@ -379,59 +350,9 @@ async fn sync_service(svc: Arc, context: Arc) -> RobotL reconcile_load_balancer(lb, svc.clone(), context).await } -/// Method to get nodes dynamically based on the pods. -/// This method will find the nodes where the target pods are deployed. -/// It will use the pod selector to find the pods and then get the nodes. -async fn get_nodes_dynamically( - svc: &Arc, - context: &Arc, -) -> RobotLBResult> { - let pod_api = kube::Api::::namespaced( - context.client.clone(), - svc.namespace() - .as_ref() - .map(String::as_str) - .unwrap_or_else(|| context.client.default_namespace()), - ); - - let Some(pod_selector) = pod_selector(svc) else { - tracing::info!( - "Service has no selector, falling back to EndpointSlice-based node discovery" - ); - return get_nodes_from_endpointslices(svc, context).await; - }; - - let label_selector = pod_selector - .iter() - .map(|(key, val)| format!("{key}={val}")) - .collect::>() - .join(","); - - let pods = pod_api - .list(&ListParams { - label_selector: Some(label_selector), - ..Default::default() - }) - .await?; - - let target_nodes = pod_nodes(&pods.items); - - let nodes_api = kube::Api::::all(context.client.clone()); - let nodes = nodes_api - .list(&ListParams::default()) - .await? - .into_iter() - .filter(|node| target_nodes.contains(&node.name_any()) && !is_excluded_from_lb(node)) - .collect::>(); - - Ok(nodes) -} - -/// Get nodes from `EndpointSlice` resources associated with a Service. -/// This method is used as a fallback when the Service has no selector, -/// such as when `EndpointSlice` resources are managed by an external controller -/// (e.g. kubevirt cloud-controller-manager). -/// It discovers target nodes by reading the `nodeName` field from each endpoint. +/// Get the nodes kube-proxy serves the service's Local traffic from. Reading the +/// slices rather than the pods also covers services whose slices another controller +/// manages (e.g. kubevirt cloud-controller-manager). async fn get_nodes_from_endpointslices( svc: &Arc, context: &Arc, @@ -448,10 +369,8 @@ async fn get_nodes_from_endpointslices( .await?; let target_nodes = eps_list - .into_iter() - .flat_map(|eps| eps.endpoints) - .filter(|ep| ep.conditions.as_ref().and_then(|c| c.ready).unwrap_or(true)) - .filter_map(|ep| ep.node_name) + .iter() + .flat_map(triggers::slice_nodes) .collect::>(); if target_nodes.is_empty() { @@ -615,7 +534,7 @@ pub async fn reconcile_load_balancer( let nodes = match node_source(&svc, context.config.dynamic_node_selector) { NodeSource::Annotation => get_nodes_by_selector(&svc, &context).await?, - NodeSource::ServiceEndpoints => get_nodes_dynamically(&svc, &context).await?, + NodeSource::ServiceEndpoints => get_nodes_from_endpointslices(&svc, &context).await?, NodeSource::AllNodes => get_all_nodes(&context).await?, }; @@ -767,8 +686,8 @@ fn error_action( mod tests { use super::{ collect_lb_services, consts, error_action, event_note, is_excluded_from_lb, - is_lb_eligible_node, is_local_traffic_policy, node_source, pod_nodes, publishes_event, - slice_triggers, success_requeue, NodeSource, + is_lb_eligible_node, is_local_traffic_policy, node_source, publishes_event, slice_triggers, + success_requeue, NodeSource, }; use k8s_openapi::{ api::core::v1::{ @@ -776,7 +695,7 @@ mod tests { }, apimachinery::pkg::apis::meta::v1::ObjectMeta, }; - use std::collections::{BTreeMap, HashSet}; + use std::collections::BTreeMap; fn service(spec: ServiceSpec) -> Service { Service { @@ -1117,46 +1036,6 @@ mod tests { assert!(!slice_triggers(&slice, None, true, &live(&slice))); } - fn pod(node: &str, phase: &str) -> k8s_openapi::api::core::v1::Pod { - k8s_openapi::api::core::v1::Pod { - spec: Some(k8s_openapi::api::core::v1::PodSpec { - node_name: Some(node.to_string()), - ..Default::default() - }), - status: Some(k8s_openapi::api::core::v1::PodStatus { - phase: Some(phase.to_string()), - pod_ip: Some("10.0.0.1".to_string()), - ..Default::default() - }), - ..Default::default() - } - } - - // An evicted pod keeps its node name, but nothing on that node serves the - // service any more, and its later deletion does not touch the endpoint slice. - #[test] - fn pods_that_finished_are_not_targets() { - let pods = [ - pod("a", "Running"), - pod("b", "Failed"), - pod("c", "Succeeded"), - pod("d", "Pending"), - ]; - assert_eq!( - pod_nodes(&pods), - HashSet::from(["a".to_string(), "d".to_string()]) - ); - } - - // A pod without an IP is not in the endpoint slice either, so its deletion - // would go unnoticed until the resync. - #[test] - fn pods_without_an_ip_are_not_targets() { - let mut scheduled = pod("a", "Pending"); - scheduled.status.as_mut().unwrap().pod_ip = None; - assert!(pod_nodes(&[scheduled]).is_empty()); - } - fn flapping_slice(ready: bool) -> k8s_openapi::api::discovery::v1::EndpointSlice { let mut slice = endpoint_slice("a"); slice.endpoints[0].conditions = Some(k8s_openapi::api::discovery::v1::EndpointConditions { @@ -1166,10 +1045,10 @@ mod tests { slice } - // Targets of a service with a selector come from its pods, so readiness of its - // endpoints moves nothing; without a selector they come from ready endpoints. + // A cordoned or not-ready node keeps running its pods, so it is the readiness of + // their endpoints that takes the node out of the targets. #[test] - fn a_readiness_flap_triggers_only_services_whose_targets_follow_readiness() { + fn a_readiness_flap_triggers_a_local_service_with_a_selector() { let with_selector = service(ServiceSpec { type_: Some("LoadBalancer".to_string()), external_traffic_policy: Some("Local".to_string()), @@ -1179,19 +1058,9 @@ mod tests { let ready = flapping_slice(true); let changes = live(&ready); slice_triggers(&ready, Some(&with_selector), true, &changes); - assert!(!slice_triggers( - &flapping_slice(false), - Some(&with_selector), - true, - &changes - )); - - let without_selector = with_policy("Local"); - let changes = live(&ready); - slice_triggers(&ready, Some(&without_selector), true, &changes); assert!(slice_triggers( &flapping_slice(false), - Some(&without_selector), + Some(&with_selector), true, &changes )); @@ -1226,39 +1095,6 @@ mod tests { } } - // A service that gains or loses its selector switches which nodes its targets - // follow, and the nodes recorded under the old rule must not hide a change. - #[test] - fn a_selector_change_does_not_hide_the_next_endpoint_change() { - let mut both = endpoint_slice("a"); - both.endpoints - .push(k8s_openapi::api::discovery::v1::Endpoint { - node_name: Some("b".to_string()), - conditions: Some(k8s_openapi::api::discovery::v1::EndpointConditions { - ready: Some(false), - ..Default::default() - }), - ..Default::default() - }); - let with_selector = service(ServiceSpec { - type_: Some("LoadBalancer".to_string()), - external_traffic_policy: Some("Local".to_string()), - selector: Some(BTreeMap::from([("app".to_string(), "web".to_string())])), - ..Default::default() - }); - let changes = live(&both); - slice_triggers(&both, Some(&with_selector), true, &changes); - // Under the new rule the nodes happen to equal what the old rule recorded. - let mut both_ready = both.clone(); - both_ready.endpoints[1].conditions = None; - assert!(slice_triggers( - &both_ready, - Some(&with_policy("Local")), - true, - &changes - )); - } - // A node whose target could not be added would otherwise wait for the resync, // since a partly successful reconcile still counts as a success. #[test] diff --git a/src/triggers.rs b/src/triggers.rs index 31d24c1..8583e74 100644 --- a/src/triggers.rs +++ b/src/triggers.rs @@ -78,15 +78,6 @@ fn fingerprint(node: &Node) -> u64 { type SliceKey = (String, String); -/// Which endpoint nodes a service's targets follow. -#[derive(Debug, Clone, Copy, PartialEq, Eq)] -pub enum Follows { - /// Targets come from the service's pods, which ignore readiness. - AllNodes, - /// Targets come from the ready endpoints. - ReadyNodes, -} - /// Tells which endpoint slice changes move balancer targets. /// /// A slice is rewritten on every readiness flip and address change, but only the @@ -104,21 +95,20 @@ pub struct SliceChanges { struct SliceState { live: HashSet, relisted: HashSet, - nodes: HashMap)>, + nodes: HashMap>, } impl SliceChanges { /// Whether the nodes of the slice changed. A slice the watch does not know to /// exist, just created or already deleted, always counts as changed. - pub fn observe(&self, slice: &EndpointSlice, follows: Follows) -> bool { + pub fn observe(&self, slice: &EndpointSlice) -> bool { let key = slice_key(slice); - let nodes = slice_nodes(slice, follows); + let nodes = slice_nodes(slice); let mut state = self.lock(); if !state.live.contains(&key) { return true; } - let recorded = (follows, nodes); - state.nodes.insert(key, recorded.clone()) != Some(recorded) + state.nodes.insert(key, nodes.clone()) != Some(nodes) } /// Drop what is known about the slice's nodes, for a service no longer observed. @@ -143,7 +133,7 @@ impl SliceChanges { Event::InitDone => { let live = std::mem::take(&mut state.relisted); let mut dropped_endpoints = false; - state.nodes.retain(|key, (_, nodes)| { + state.nodes.retain(|key, nodes| { let kept = live.contains(key); dropped_endpoints |= !kept && !nodes.is_empty(); kept @@ -161,7 +151,7 @@ impl SliceChanges { state .nodes .remove(&key) - .is_some_and(|(_, nodes)| !nodes.is_empty()) + .is_some_and(|nodes| !nodes.is_empty()) } } } @@ -177,14 +167,23 @@ fn slice_key(slice: &K) -> SliceKey { (slice.namespace().unwrap_or_default(), slice.name_any()) } -fn slice_nodes(slice: &EndpointSlice, follows: Follows) -> BTreeSet { +/// The nodes kube-proxy serves the slice's Local traffic from. +/// +/// These are the nodes with a ready endpoint, or with a terminating one that still +/// serves, which kube-proxy falls back to while a node has no ready endpoint. +#[must_use] +pub fn slice_nodes(slice: &EndpointSlice) -> BTreeSet { slice .endpoints .iter() - // The API treats a missing ready condition as ready. .filter(|endpoint| { - follows == Follows::AllNodes - || endpoint.conditions.as_ref().and_then(|c| c.ready) != Some(false) + let conditions = endpoint.conditions.as_ref(); + // kube-proxy reads a missing ready or serving condition as true, and a + // missing terminating one as false. + let ready = conditions.and_then(|c| c.ready) != Some(false); + let serving = conditions.and_then(|c| c.serving) != Some(false); + let terminating = conditions.and_then(|c| c.terminating) == Some(true); + ready || (serving && terminating) }) .filter_map(|endpoint| endpoint.node_name.clone()) .collect() @@ -199,7 +198,7 @@ pub fn slice_service(slice: &EndpointSlice) -> Option> { #[cfg(test)] mod tests { - use super::{slice_service, Follows, NodeChanges, SliceChanges}; + use super::{slice_nodes, slice_service, NodeChanges, SliceChanges}; use k8s_openapi::{ api::{ core::v1::{Node, NodeAddress, NodeCondition, NodeSpec, NodeStatus}, @@ -361,77 +360,128 @@ mod tests { } } - fn tracked(initial: &EndpointSlice, follows: Follows) -> SliceChanges { + fn tracked(initial: &EndpointSlice) -> SliceChanges { let changes = SliceChanges::default(); assert!(!changes.track(&Event::::Init)); assert!(!changes.track(&Event::InitApply(initial.clone()))); assert!(!changes.track(&Event::::InitDone)); - changes.observe(initial, follows); + changes.observe(initial); changes } #[test] fn a_slice_created_after_the_listing_is_tracked_like_a_listed_one() { - let changes = tracked(&slice(&[]), Follows::ReadyNodes); + let changes = tracked(&slice(&[])); let created = slice(&[("10.0.0.1", "a", true)]); let mut created = created; created.metadata.name = Some("web-fghij".to_string()); assert!(!changes.track(&Event::Apply(created.clone()))); - assert!(changes.observe(&created, Follows::ReadyNodes)); - assert!(!changes.observe(&created, Follows::ReadyNodes)); + assert!(changes.observe(&created)); + assert!(!changes.observe(&created)); assert!(changes.track(&Event::Delete(created))); } #[test] fn a_slice_not_known_to_exist_triggers_a_reconcile() { let changes = SliceChanges::default(); - assert!(changes.observe(&slice(&[("10.0.0.1", "a", true)]), Follows::ReadyNodes)); + assert!(changes.observe(&slice(&[("10.0.0.1", "a", true)]))); } // Replicas flapping readiness or changing addresses on the same nodes rewrite the // slice all the time, and none of that moves a balancer target. #[test] fn endpoint_churn_on_the_same_nodes_triggers_nothing() { - for follows in [Follows::AllNodes, Follows::ReadyNodes] { - let changes = tracked( - &slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "a", true)]), - follows, - ); - let moved = slice(&[("10.0.0.3", "a", true), ("10.0.0.2", "a", true)]); - assert!(!changes.observe(&moved, follows)); - let flapped = slice(&[("10.0.0.3", "a", false), ("10.0.0.2", "a", true)]); - assert!(!changes.observe(&flapped, follows)); - } + let changes = tracked(&slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "a", true)])); + let moved = slice(&[("10.0.0.3", "a", true), ("10.0.0.2", "a", true)]); + assert!(!changes.observe(&moved)); + let flapped = slice(&[("10.0.0.3", "a", false), ("10.0.0.2", "a", true)]); + assert!(!changes.observe(&flapped)); } #[test] fn endpoint_targets_follow_the_last_ready_endpoint_of_a_node() { - let changes = tracked(&slice(&[("10.0.0.1", "a", true)]), Follows::ReadyNodes); + let changes = tracked(&slice(&[("10.0.0.1", "a", true)])); let added = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", true)]); - assert!(changes.observe(&added, Follows::ReadyNodes)); + assert!(changes.observe(&added)); let unready = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", false)]); - assert!(changes.observe(&unready, Follows::ReadyNodes)); + assert!(changes.observe(&unready)); + } + + fn conditions( + ready: Option, + serving: Option, + terminating: Option, + ) -> EndpointConditions { + EndpointConditions { + ready, + serving, + terminating, + } } - // Pod-based targets ignore readiness, so a flap costs nothing there, while an - // endpoint leaving the slice, terminating or not, moves a target. + // kube-proxy serves Local traffic from a node's ready endpoints, or from its + // terminating ones that still serve when none is ready, and reads a missing + // ready or serving condition as true and a missing terminating one as false: + // https://github.com/kubernetes/kubernetes/blob/d7c57fb776cbf2554d352a835459194ee62cf751/pkg/proxy/endpointslicecache.go#L209-L211 #[test] - fn pod_targets_ignore_readiness_and_follow_endpoints_leaving() { - let both = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", true)]); - let changes = tracked(&both, Follows::AllNodes); - let unready = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", false)]); - assert!(!changes.observe(&unready, Follows::AllNodes)); - assert!(changes.observe(&slice(&[("10.0.0.1", "a", true)]), Follows::AllNodes)); + fn nodes_kube_proxy_serves_local_traffic_from_are_targets() { + let mut listed = slice(&[]); + for (node, conditions) in [ + ("ready", conditions(Some(true), Some(true), Some(false))), + ( + "not-ready", + conditions(Some(false), Some(false), Some(false)), + ), + ("draining", conditions(Some(false), Some(true), Some(true))), + ("drained", conditions(Some(false), Some(false), Some(true))), + ( + "draining-unknown-serving", + conditions(Some(false), None, Some(true)), + ), + ( + "not-ready-unknown-terminating", + conditions(Some(false), Some(true), None), + ), + ("unknown-ready", conditions(None, None, None)), + ] { + listed.endpoints.push(Endpoint { + addresses: vec!["10.0.0.1".to_string()], + node_name: Some(node.to_string()), + conditions: Some(conditions), + ..Default::default() + }); + } + listed.endpoints.push(Endpoint { + addresses: vec!["10.0.0.2".to_string()], + node_name: Some("no-conditions".to_string()), + ..Default::default() + }); + listed.endpoints.push(Endpoint { + addresses: vec!["10.0.0.3".to_string()], + ..Default::default() + }); + assert_eq!( + slice_nodes(&listed), + [ + "ready", + "draining", + "draining-unknown-serving", + "unknown-ready", + "no-conditions" + ] + .map(String::from) + .into() + ); } #[test] fn deleting_a_slice_with_endpoints_triggers_a_reconcile_of_everything() { let initial = slice(&[("10.0.0.1", "a", true)]); - let changes = tracked(&initial, Follows::ReadyNodes); + let changes = tracked(&initial); assert!(changes.track(&Event::Delete(initial.clone()))); // The mapper may see the deletion after the watch did; the slice is gone by // then, so it reconciles the service rather than recording the slice again. - assert!(changes.observe(&initial, Follows::ReadyNodes)); + assert!(changes.observe(&initial)); assert!(!changes.track(&Event::Delete(initial))); } @@ -440,10 +490,10 @@ mod tests { #[test] fn a_relisting_forgets_slices_that_are_gone() { let initial = slice(&[("10.0.0.1", "a", true)]); - let changes = tracked(&initial, Follows::ReadyNodes); + let changes = tracked(&initial); assert!(!changes.track(&Event::::Init)); assert!(changes.track(&Event::::InitDone)); - assert!(changes.observe(&initial, Follows::ReadyNodes)); + assert!(changes.observe(&initial)); } // A service that stops using the Local policy stops being observed, and its old @@ -451,15 +501,15 @@ mod tests { #[test] fn a_forgotten_slice_counts_as_changed() { let initial = slice(&[("10.0.0.1", "a", true)]); - let changes = tracked(&initial, Follows::ReadyNodes); + let changes = tracked(&initial); changes.forget(&initial); - assert!(changes.observe(&initial, Follows::ReadyNodes)); + assert!(changes.observe(&initial)); } #[test] fn a_relisting_that_drops_only_empty_slices_triggers_nothing() { let empty = slice(&[]); - let changes = tracked(&empty, Follows::ReadyNodes); + let changes = tracked(&empty); assert!(!changes.track(&Event::::Init)); assert!(!changes.track(&Event::::InitDone)); }