diff --git a/README.md b/README.md index 02c9aae..47f69c2 100644 --- a/README.md +++ b/README.md @@ -35,10 +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 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. @@ -70,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 @@ -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 @@ -142,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 7c3a25f..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, @@ -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..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, @@ -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,142 @@ 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(()) +} + +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 those 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; + } + changes.observe(slice) +} + +/// 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)] @@ -249,68 +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) = svc - .spec - .as_ref() - .and_then(|spec| spec.selector.clone()) - .filter(|s| !s.is_empty()) - 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 = pods - .iter() - .map(|pod| pod.spec.clone().unwrap_or_default().node_name) - .flatten() - .collect::>(); - - 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, @@ -327,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() { @@ -494,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?, }; @@ -528,10 +568,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 +614,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 +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, publishes_event, NodeSource, + is_lb_eligible_node, is_local_traffic_policy, node_source, publishes_event, slice_triggers, + success_requeue, NodeSource, }; use k8s_openapi::{ api::core::v1::{ @@ -908,4 +964,144 @@ 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 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 + } + + // 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_a_local_service_with_a_selector() { + 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 + )); + } + + // 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 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..8583e74 --- /dev/null +++ b/src/triggers.rs @@ -0,0 +1,516 @@ +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); + +/// 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) -> bool { + let key = slice_key(slice); + let nodes = slice_nodes(slice); + let mut state = self.lock(); + if !state.live.contains(&key) { + return true; + } + state.nodes.insert(key, nodes.clone()) != Some(nodes) + } + + /// 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()) +} + +/// 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() + .filter(|endpoint| { + 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() +} + +#[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_nodes, slice_service, 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) -> 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); + changes + } + + #[test] + fn a_slice_created_after_the_listing_is_tracked_like_a_listed_one() { + 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)); + 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)]))); + } + + // 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() { + 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)])); + let added = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", true)]); + assert!(changes.observe(&added)); + let unready = slice(&[("10.0.0.1", "a", true), ("10.0.0.2", "b", false)]); + assert!(changes.observe(&unready)); + } + + fn conditions( + ready: Option, + serving: Option, + terminating: Option, + ) -> EndpointConditions { + EndpointConditions { + ready, + serving, + terminating, + } + } + + // 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 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); + 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)); + 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); + assert!(!changes.track(&Event::::Init)); + assert!(changes.track(&Event::::InitDone)); + assert!(changes.observe(&initial)); + } + + // 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); + changes.forget(&initial); + assert!(changes.observe(&initial)); + } + + #[test] + fn a_relisting_that_drops_only_empty_slices_triggers_nothing() { + let empty = slice(&[]); + let changes = tracked(&empty); + assert!(!changes.track(&Event::::Init)); + assert!(!changes.track(&Event::::InitDone)); + } +}