Skip to content
Draft
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
218 changes: 74 additions & 144 deletions hcloud/load_balancers.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,18 @@ import (
"context"
"errors"
"fmt"
"slices"

corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/labels"
cloudprovider "k8s.io/cloud-provider"
"k8s.io/klog/v2"

"github.com/hetznercloud/hcloud-cloud-controller-manager/internal/annotation"
"github.com/hetznercloud/hcloud-cloud-controller-manager/internal/config"
"github.com/hetznercloud/hcloud-cloud-controller-manager/internal/hcops"
"github.com/hetznercloud/hcloud-cloud-controller-manager/internal/lbspec"
"github.com/hetznercloud/hcloud-cloud-controller-manager/internal/metrics"
"github.com/hetznercloud/hcloud-go/v2/hcloud"
"github.com/hetznercloud/hcloud-go/v2/hcloud/exp/kit/sliceutil"
)

// LoadBalancerOps defines the Load Balancer related operations required by
Expand All @@ -23,7 +24,7 @@ type LoadBalancerOps interface {
GetByName(ctx context.Context, name string) (*hcloud.LoadBalancer, error)
GetByID(ctx context.Context, id int64) (*hcloud.LoadBalancer, error)
GetByK8SServiceUID(ctx context.Context, svc *corev1.Service) (*hcloud.LoadBalancer, error)
Create(ctx context.Context, lbName string, service *corev1.Service) (*hcloud.LoadBalancer, error)
Create(ctx context.Context, service *corev1.Service) (*hcloud.LoadBalancer, error)
Delete(ctx context.Context, lb *hcloud.LoadBalancer) error
ReconcileHCLB(ctx context.Context, lb *hcloud.LoadBalancer, svc *corev1.Service) (bool, error)
ReconcileHCLBTargets(ctx context.Context, lb *hcloud.LoadBalancer, svc *corev1.Service, nodes []*corev1.Node) (bool, error)
Expand All @@ -42,33 +43,18 @@ func newLoadBalancers(lbOps LoadBalancerOps, lbCfg *config.LoadBalancerConfigura
}
}

func matchNodeSelector(svc *corev1.Service, nodes []*corev1.Node) ([]*corev1.Node, error) {
var selectedNodes []*corev1.Node

selector := labels.Everything()
if v, err := annotation.LBNodeSelector.FromService(svc); err == nil {
parsed, err := labels.Parse(v)
if err != nil {
return nil, fmt.Errorf("unable to parse the node-selector annotation: %w", err)
}
selector = parsed
}

for _, n := range nodes {
if selector.Matches(labels.Set(n.GetLabels())) {
selectedNodes = append(selectedNodes, n)
}
}

return selectedNodes, nil
}

func (l *loadBalancers) GetLoadBalancer(
ctx context.Context, _ string, service *corev1.Service,
) (status *corev1.LoadBalancerStatus, exists bool, err error) {
const op = "hcloud/loadBalancers.GetLoadBalancer"
metrics.OperationCalled.WithLabelValues(op).Inc()

// The lookup comes before resolving the annotations on purpose. The service
// controller calls us before deleting a Load Balancer and only removes its
// finalizer once we returned without an error, so reporting that a Load
// Balancer does not exist must not depend on the annotations being valid.
// Otherwise a Service that never got a Load Balancer because one of its
// annotations is invalid could not be deleted either.
lb, err := l.lbOps.GetByK8SServiceUID(ctx, service)
if err != nil {
if errors.Is(err, hcops.ErrNotFound) {
Expand All @@ -77,56 +63,41 @@ func (l *loadBalancers) GetLoadBalancer(
return nil, false, fmt.Errorf("%s: %w", op, err)
}

if v, err := annotation.LBHostname.FromService(service); err == nil {
return &corev1.LoadBalancerStatus{
Ingress: []corev1.LoadBalancerIngress{{Hostname: v}},
}, true, nil
}

ingress, err := l.buildLoadBalancerStatusIngress(lb, service)
spec, err := lbspec.Resolve(service, *l.cfg)
if err != nil {
return nil, false, fmt.Errorf("%s: %w", op, err)
}

return &corev1.LoadBalancerStatus{Ingress: ingress}, true, nil
return &corev1.LoadBalancerStatus{Ingress: l.buildLoadBalancerStatusIngress(lb, spec)}, true, nil
}

func (l *loadBalancers) GetLoadBalancerName(_ context.Context, _ string, service *corev1.Service) string {
if v, err := annotation.LBName.FromService(service); err == nil {
return v
}
return cloudprovider.DefaultLoadBalancerName(service)
return lbspec.Name(service)
}

func (l *loadBalancers) EnsureLoadBalancer(
ctx context.Context, clusterName string, svc *corev1.Service, nodes []*corev1.Node,
ctx context.Context, _ string, svc *corev1.Service, nodes []*corev1.Node,
) (*corev1.LoadBalancerStatus, error) {
const op = "hcloud/loadBalancers.EnsureLoadBalancer"
metrics.OperationCalled.WithLabelValues(op).Inc()

var (
reload bool
lb *hcloud.LoadBalancer
err error
selectedNodes []*corev1.Node
reload bool
lb *hcloud.LoadBalancer
err error
)

selectedNodes, err = matchNodeSelector(svc, nodes)
spec, err := lbspec.Resolve(svc, *l.cfg)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}

nodeNames := make([]string, len(selectedNodes))
for i, n := range selectedNodes {
nodeNames[i] = n.Name
}
selectedNodes := filterNodes(spec.NodeSelector, nodes)
nodeNames := sliceutil.Transform(selectedNodes, func(n *corev1.Node) string {
return n.GetName()
})
klog.InfoS("ensure Load Balancer", "op", op, "service", svc.Name, "nodes", nodeNames)

lb, err = l.lbOps.GetByK8SServiceUID(ctx, svc)
if err != nil && !errors.Is(err, hcops.ErrNotFound) {
return nil, fmt.Errorf("%s: %w", op, err)
}

// Try the load balancer's name if we were not able to find it using the
// service UID. This is required for two reasons:
//
Expand All @@ -135,20 +106,23 @@ func (l *loadBalancers) EnsureLoadBalancer(
//
// 2. Import of load balancers which were created by other means but
// should be re-used by the cloud controller manager.
lbName := l.GetLoadBalancerName(ctx, clusterName, svc)
lb, err = l.lbOps.GetByK8SServiceUID(ctx, svc)
if err != nil && !errors.Is(err, hcops.ErrNotFound) {
return nil, fmt.Errorf("%s: %w", op, err)
}

// Not found by UID label; try name
if errors.Is(err, hcops.ErrNotFound) {
lb, err = l.lbOps.GetByName(ctx, lbName)
if err != nil && !errors.Is(err, hcops.ErrNotFound) {
return nil, fmt.Errorf("%s: %w", op, err)
}
lb, err = l.lbOps.GetByName(ctx, spec.Name)
}

// If we were still not able to find the load balancer we create it.
// New Load Balancer -> create it
if errors.Is(err, hcops.ErrNotFound) {
lb, err = l.lbOps.Create(ctx, lbName, svc)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}
lb, err = l.lbOps.Create(ctx, svc)
}

if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}

lbChanged, err := l.lbOps.ReconcileHCLB(ctx, lb, svc)
Expand Down Expand Up @@ -192,62 +166,42 @@ func (l *loadBalancers) EnsureLoadBalancer(
}
}

// Either set the Hostname or the IPs (below).
// See: https://github.com/kubernetes/kubernetes/issues/66607
if v, err := annotation.LBHostname.FromService(svc); err == nil {
return &corev1.LoadBalancerStatus{
Ingress: []corev1.LoadBalancerIngress{{Hostname: v}},
}, nil
}

ingress, err := l.buildLoadBalancerStatusIngress(lb, svc)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}

return &corev1.LoadBalancerStatus{Ingress: ingress}, nil
return &corev1.LoadBalancerStatus{Ingress: l.buildLoadBalancerStatusIngress(lb, spec)}, nil
}

func (l *loadBalancers) buildLoadBalancerStatusIngress(lb *hcloud.LoadBalancer, svc *corev1.Service) ([]corev1.LoadBalancerIngress, error) {
// buildLoadBalancerStatusIngress reports the addresses the Service is reachable
// on. A configured hostname replaces the IPs entirely.
// See: https://github.com/kubernetes/kubernetes/issues/66607
func (l *loadBalancers) buildLoadBalancerStatusIngress(lb *hcloud.LoadBalancer, spec lbspec.Spec) []corev1.LoadBalancerIngress {
const op = "hcloud/loadBalancers.getLoadBalancerStatusIngress"
metrics.OperationCalled.WithLabelValues(op).Inc()

var ingress []corev1.LoadBalancerIngress
ipMode := corev1.LoadBalancerIPModeVIP

proxyProtocolEnabled, err := l.getProxyProtocolEnabled(svc)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
if spec.Hostname != "" {
return []corev1.LoadBalancerIngress{{Hostname: spec.Hostname}}
}

if proxyProtocolEnabled {
ipMode := corev1.LoadBalancerIPModeVIP
if spec.Service.ProxyProtocol != nil && *spec.Service.ProxyProtocol {
ipMode = corev1.LoadBalancerIPModeProxy
}

var ingress []corev1.LoadBalancerIngress

if lb.PublicNet.Enabled {
ingress = append(ingress, corev1.LoadBalancerIngress{
IP: lb.PublicNet.IPv4.IP.String(),
IPMode: &ipMode,
})

ipv6Enabled, err := l.getIPv6Enabled(svc)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}
if ipv6Enabled {
if spec.IPv6 {
ingress = append(ingress, corev1.LoadBalancerIngress{
IP: lb.PublicNet.IPv6.IP.String(),
IPMode: &ipMode,
})
}
}

privateIngressEnabled, err := l.getPrivateIngressEnabled(svc)
if err != nil {
return nil, fmt.Errorf("%s: %w", op, err)
}

if privateIngressEnabled {
if spec.PrivateIngress {
for _, privateNet := range lb.PrivateNet {
ingress = append(ingress, corev1.LoadBalancerIngress{
IP: privateNet.IP.String(),
Expand All @@ -256,62 +210,30 @@ func (l *loadBalancers) buildLoadBalancerStatusIngress(lb *hcloud.LoadBalancer,
}
}

return ingress, nil
}

func (l *loadBalancers) getPrivateIngressEnabled(svc *corev1.Service) (bool, error) {
disable, err := annotation.LBDisablePrivateIngress.FromService(svc)
if err == nil {
return !disable, nil
}
if errors.Is(err, annotation.ErrNotSet) {
return l.cfg.PrivateIngressEnabled, nil
}
return true, err
}

func (l *loadBalancers) getProxyProtocolEnabled(svc *corev1.Service) (bool, error) {
enable, err := annotation.LBSvcProxyProtocol.FromService(svc)
if err == nil {
return enable, nil
}
if errors.Is(err, annotation.ErrNotSet) {
if l.cfg.ProxyProtocolEnabled == nil {
return false, nil
}
return *l.cfg.ProxyProtocolEnabled, nil
}
return false, err
}

func (l *loadBalancers) getIPv6Enabled(svc *corev1.Service) (bool, error) {
disable, err := annotation.LBIPv6Disabled.FromService(svc)
if err == nil {
return !disable, nil
}
if errors.Is(err, annotation.ErrNotSet) {
return l.cfg.IPv6Enabled, nil
}
return true, err
return ingress
}

func (l *loadBalancers) UpdateLoadBalancer(
ctx context.Context, clusterName string, svc *corev1.Service, nodes []*corev1.Node,
ctx context.Context,
_ string,
svc *corev1.Service,
nodes []*corev1.Node,
) error {
const op = "hcloud/loadBalancers.UpdateLoadBalancer"
metrics.OperationCalled.WithLabelValues(op).Inc()

var (
lb *hcloud.LoadBalancer
err error
selectedNodes []*corev1.Node
lb *hcloud.LoadBalancer
err error
)

selectedNodes, err = matchNodeSelector(svc, nodes)
spec, err := lbspec.Resolve(svc, *l.cfg)
if err != nil {
return fmt.Errorf("%s: %w", op, err)
}

selectedNodes := filterNodes(spec.NodeSelector, nodes)

nodeNames := make([]string, len(selectedNodes))
for i, n := range selectedNodes {
nodeNames[i] = n.Name
Expand All @@ -320,15 +242,13 @@ func (l *loadBalancers) UpdateLoadBalancer(

lb, err = l.lbOps.GetByK8SServiceUID(ctx, svc)
if errors.Is(err, hcops.ErrNotFound) {
lbName := l.GetLoadBalancerName(ctx, clusterName, svc)

lb, err = l.lbOps.GetByName(ctx, lbName)
if errors.Is(err, hcops.ErrNotFound) {
return nil
}
// further error types handled below
lb, err = l.lbOps.GetByName(ctx, spec.Name)
}
if err != nil {
switch {
case errors.Is(err, hcops.ErrNotFound):
// Nothing to do, the Load Balancer does not exist.
return nil
case err != nil:
return fmt.Errorf("%s: %w", op, err)
}

Expand Down Expand Up @@ -372,3 +292,13 @@ func (l *loadBalancers) EnsureLoadBalancerDeleted(ctx context.Context, _ string,

return nil
}

func filterNodes(selector labels.Selector, nodes []*corev1.Node) []*corev1.Node {
if selector.Empty() {
return nodes
}

return slices.DeleteFunc(slices.Clone(nodes), func(n *corev1.Node) bool {
return !selector.Matches(labels.Set(n.GetLabels()))
})
}
Loading
Loading