Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -72,7 +72,7 @@ Options:
--default-network <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_LB_RETRIES>
Default load balancer healthcheck retries cound [env: ROBOTLB_DEFAULT_LB_RETRIES=] [default: 3]
--default-lb-timeout <DEFAULT_LB_TIMEOUT>
Expand Down Expand Up @@ -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:
Expand Down
2 changes: 1 addition & 1 deletion helm/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
2 changes: 1 addition & 1 deletion src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ pub struct OperatorConfig {
#[arg(long, env = "ROBOTLB_DEFAULT_NETWORK", default_value = None)]
pub default_network: Option<String>,

/// 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,
Expand Down
196 changes: 16 additions & 180 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<String> {
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<std::collections::BTreeMap<String, String>> {
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")
Expand All @@ -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>,
Expand All @@ -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
Expand Down Expand Up @@ -379,59 +350,9 @@ async fn sync_service(svc: Arc<Service>, context: Arc<CurrentContext>) -> 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<Service>,
context: &Arc<CurrentContext>,
) -> RobotLBResult<Vec<Node>> {
let pod_api = kube::Api::<Pod>::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::<Vec<_>>()
.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::<Node>::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::<Vec<_>>();

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<Service>,
context: &Arc<CurrentContext>,
Expand All @@ -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::<HashSet<_>>();

if target_nodes.is_empty() {
Expand Down Expand Up @@ -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?,
};

Expand Down Expand Up @@ -767,16 +686,16 @@ 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::{
Node, NodeCondition, NodeSpec, NodeStatus, Service, ServicePort, ServiceSpec,
},
apimachinery::pkg::apis::meta::v1::ObjectMeta,
};
use std::collections::{BTreeMap, HashSet};
use std::collections::BTreeMap;

fn service(spec: ServiceSpec) -> Service {
Service {
Expand Down Expand Up @@ -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 {
Expand All @@ -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()),
Expand All @@ -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
));
Expand Down Expand Up @@ -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]
Expand Down
Loading
Loading