|
|
@ -12,11 +12,13 @@ import ( |
|
|
|
"github.com/k3s-io/k3s/pkg/version" |
|
|
|
"github.com/k3s-io/k3s/pkg/version" |
|
|
|
"github.com/rancher/wrangler/pkg/condition" |
|
|
|
"github.com/rancher/wrangler/pkg/condition" |
|
|
|
coreclient "github.com/rancher/wrangler/pkg/generated/controllers/core/v1" |
|
|
|
coreclient "github.com/rancher/wrangler/pkg/generated/controllers/core/v1" |
|
|
|
|
|
|
|
discoveryclient "github.com/rancher/wrangler/pkg/generated/controllers/discovery/v1" |
|
|
|
"github.com/rancher/wrangler/pkg/merr" |
|
|
|
"github.com/rancher/wrangler/pkg/merr" |
|
|
|
"github.com/rancher/wrangler/pkg/objectset" |
|
|
|
"github.com/rancher/wrangler/pkg/objectset" |
|
|
|
"github.com/sirupsen/logrus" |
|
|
|
"github.com/sirupsen/logrus" |
|
|
|
apps "k8s.io/api/apps/v1" |
|
|
|
apps "k8s.io/api/apps/v1" |
|
|
|
core "k8s.io/api/core/v1" |
|
|
|
core "k8s.io/api/core/v1" |
|
|
|
|
|
|
|
discovery "k8s.io/api/discovery/v1" |
|
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors" |
|
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors" |
|
|
|
meta "k8s.io/apimachinery/pkg/apis/meta/v1" |
|
|
|
meta "k8s.io/apimachinery/pkg/apis/meta/v1" |
|
|
|
"k8s.io/apimachinery/pkg/labels" |
|
|
|
"k8s.io/apimachinery/pkg/labels" |
|
|
@ -48,9 +50,11 @@ const ( |
|
|
|
func (k *k3s) Register(ctx context.Context, |
|
|
|
func (k *k3s) Register(ctx context.Context, |
|
|
|
nodes coreclient.NodeController, |
|
|
|
nodes coreclient.NodeController, |
|
|
|
pods coreclient.PodController, |
|
|
|
pods coreclient.PodController, |
|
|
|
|
|
|
|
endpointslices discoveryclient.EndpointSliceController, |
|
|
|
) error { |
|
|
|
) error { |
|
|
|
nodes.OnChange(ctx, controllerName, k.onChangeNode) |
|
|
|
nodes.OnChange(ctx, controllerName, k.onChangeNode) |
|
|
|
pods.OnChange(ctx, controllerName, k.onChangePod) |
|
|
|
pods.OnChange(ctx, controllerName, k.onChangePod) |
|
|
|
|
|
|
|
endpointslices.OnChange(ctx, controllerName, k.onChangeEndpointSlice) |
|
|
|
|
|
|
|
|
|
|
|
if err := k.createServiceLBNamespace(ctx); err != nil { |
|
|
|
if err := k.createServiceLBNamespace(ctx); err != nil { |
|
|
|
return err |
|
|
|
return err |
|
|
@ -135,6 +139,22 @@ func (k *k3s) onChangeNode(key string, node *core.Node) (*core.Node, error) { |
|
|
|
return node, nil |
|
|
|
return node, nil |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// onChangeEndpointSlice handles changes to EndpointSlices. This is used to ensure that LoadBalancer
|
|
|
|
|
|
|
|
// addresses only list Nodes with ready Pods, when their ExternalTrafficPolicy is set to Local.
|
|
|
|
|
|
|
|
func (k *k3s) onChangeEndpointSlice(key string, eps *discovery.EndpointSlice) (*discovery.EndpointSlice, error) { |
|
|
|
|
|
|
|
if eps == nil { |
|
|
|
|
|
|
|
return nil, nil |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
serviceName, ok := eps.Labels[discovery.LabelServiceName] |
|
|
|
|
|
|
|
if !ok { |
|
|
|
|
|
|
|
return eps, nil |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
k.workqueue.Add(eps.Namespace + "/" + serviceName) |
|
|
|
|
|
|
|
return eps, nil |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// runWorker dequeues Service changes from the work queue
|
|
|
|
// runWorker dequeues Service changes from the work queue
|
|
|
|
// We run a lightweight work queue to handle service updates. We don't need the full overhead
|
|
|
|
// We run a lightweight work queue to handle service updates. We don't need the full overhead
|
|
|
|
// of a wrangler service controller and shared informer cache, but we do want to run changes
|
|
|
|
// of a wrangler service controller and shared informer cache, but we do want to run changes
|
|
|
@ -219,16 +239,37 @@ func (k *k3s) getDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
// getStatus returns a LoadBalancerStatus listing ingress IPs for all ready pods
|
|
|
|
// getStatus returns a LoadBalancerStatus listing ingress IPs for all ready pods
|
|
|
|
// matching the selected service.
|
|
|
|
// matching the selected service.
|
|
|
|
func (k *k3s) getStatus(svc *core.Service) (*core.LoadBalancerStatus, error) { |
|
|
|
func (k *k3s) getStatus(svc *core.Service) (*core.LoadBalancerStatus, error) { |
|
|
|
pods, err := k.podCache.List(k.LBNamespace, labels.SelectorFromSet(map[string]string{ |
|
|
|
var readyNodes map[string]bool |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if servicehelper.RequestsOnlyLocalTraffic(svc) { |
|
|
|
|
|
|
|
readyNodes = map[string]bool{} |
|
|
|
|
|
|
|
eps, err := k.endpointsCache.List(svc.Namespace, labels.SelectorFromSet(labels.Set{ |
|
|
|
|
|
|
|
discovery.LabelServiceName: svc.Name, |
|
|
|
|
|
|
|
})) |
|
|
|
|
|
|
|
if err != nil { |
|
|
|
|
|
|
|
return nil, err |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for _, ep := range eps { |
|
|
|
|
|
|
|
for _, endpoint := range ep.Endpoints { |
|
|
|
|
|
|
|
isPod := endpoint.TargetRef != nil && endpoint.TargetRef.Kind == "Pod" |
|
|
|
|
|
|
|
isReady := endpoint.Conditions.Ready != nil && *endpoint.Conditions.Ready |
|
|
|
|
|
|
|
if isPod && isReady && endpoint.NodeName != nil { |
|
|
|
|
|
|
|
readyNodes[*endpoint.NodeName] = true |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
pods, err := k.podCache.List(k.LBNamespace, labels.SelectorFromSet(labels.Set{ |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
|
})) |
|
|
|
})) |
|
|
|
|
|
|
|
|
|
|
|
if err != nil { |
|
|
|
if err != nil { |
|
|
|
return nil, err |
|
|
|
return nil, err |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
expectedIPs, err := k.podIPs(pods, svc) |
|
|
|
expectedIPs, err := k.podIPs(pods, svc, readyNodes) |
|
|
|
if err != nil { |
|
|
|
if err != nil { |
|
|
|
return nil, err |
|
|
|
return nil, err |
|
|
|
} |
|
|
|
} |
|
|
@ -267,7 +308,7 @@ func (k *k3s) patchStatus(svc *core.Service, previousStatus, newStatus *core.Loa |
|
|
|
// podIPs returns a list of IPs for Nodes hosting ServiceLB Pods.
|
|
|
|
// podIPs returns a list of IPs for Nodes hosting ServiceLB Pods.
|
|
|
|
// If at least one node has External IPs available, only external IPs are returned.
|
|
|
|
// If at least one node has External IPs available, only external IPs are returned.
|
|
|
|
// If no nodes have External IPs set, the Internal IPs of all nodes running pods are returned.
|
|
|
|
// If no nodes have External IPs set, the Internal IPs of all nodes running pods are returned.
|
|
|
|
func (k *k3s) podIPs(pods []*core.Pod, svc *core.Service) ([]string, error) { |
|
|
|
func (k *k3s) podIPs(pods []*core.Pod, svc *core.Service, readyNodes map[string]bool) ([]string, error) { |
|
|
|
// Go doesn't have sets so we stuff things into a map of bools and then get lists of keys
|
|
|
|
// Go doesn't have sets so we stuff things into a map of bools and then get lists of keys
|
|
|
|
// to determine the unique set of IPs in use by pods.
|
|
|
|
// to determine the unique set of IPs in use by pods.
|
|
|
|
extIPs := map[string]bool{} |
|
|
|
extIPs := map[string]bool{} |
|
|
@ -280,6 +321,9 @@ func (k *k3s) podIPs(pods []*core.Pod, svc *core.Service) ([]string, error) { |
|
|
|
if !Ready.IsTrue(pod) { |
|
|
|
if !Ready.IsTrue(pod) { |
|
|
|
continue |
|
|
|
continue |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
if readyNodes != nil && !readyNodes[pod.Spec.NodeName] { |
|
|
|
|
|
|
|
continue |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
node, err := k.nodeCache.Get(pod.Spec.NodeName) |
|
|
|
node, err := k.nodeCache.Get(pod.Spec.NodeName) |
|
|
|
if apierrors.IsNotFound(err) { |
|
|
|
if apierrors.IsNotFound(err) { |
|
|
@ -405,17 +449,27 @@ func (k *k3s) deleteDaemonSet(ctx context.Context, svc *core.Service) error { |
|
|
|
func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
name := generateName(svc) |
|
|
|
name := generateName(svc) |
|
|
|
oneInt := intstr.FromInt(1) |
|
|
|
oneInt := intstr.FromInt(1) |
|
|
|
|
|
|
|
localTraffic := servicehelper.RequestsOnlyLocalTraffic(svc) |
|
|
|
sourceRanges, err := servicehelper.GetLoadBalancerSourceRanges(svc) |
|
|
|
sourceRanges, err := servicehelper.GetLoadBalancerSourceRanges(svc) |
|
|
|
if err != nil { |
|
|
|
if err != nil { |
|
|
|
return nil, err |
|
|
|
return nil, err |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
var sysctls []core.Sysctl |
|
|
|
|
|
|
|
for _, ipFamily := range svc.Spec.IPFamilies { |
|
|
|
|
|
|
|
switch ipFamily { |
|
|
|
|
|
|
|
case core.IPv4Protocol: |
|
|
|
|
|
|
|
sysctls = append(sysctls, core.Sysctl{Name: "net.ipv4.ip_forward", Value: "1"}) |
|
|
|
|
|
|
|
case core.IPv6Protocol: |
|
|
|
|
|
|
|
sysctls = append(sysctls, core.Sysctl{Name: "net.ipv6.conf.all.forwarding", Value: "1"}) |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
ds := &apps.DaemonSet{ |
|
|
|
ds := &apps.DaemonSet{ |
|
|
|
ObjectMeta: meta.ObjectMeta{ |
|
|
|
ObjectMeta: meta.ObjectMeta{ |
|
|
|
Name: name, |
|
|
|
Name: name, |
|
|
|
Namespace: k.LBNamespace, |
|
|
|
Namespace: k.LBNamespace, |
|
|
|
Labels: map[string]string{ |
|
|
|
Labels: labels.Set{ |
|
|
|
nodeSelectorLabel: "false", |
|
|
|
nodeSelectorLabel: "false", |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
@ -427,13 +481,13 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
}, |
|
|
|
}, |
|
|
|
Spec: apps.DaemonSetSpec{ |
|
|
|
Spec: apps.DaemonSetSpec{ |
|
|
|
Selector: &meta.LabelSelector{ |
|
|
|
Selector: &meta.LabelSelector{ |
|
|
|
MatchLabels: map[string]string{ |
|
|
|
MatchLabels: labels.Set{ |
|
|
|
"app": name, |
|
|
|
"app": name, |
|
|
|
}, |
|
|
|
}, |
|
|
|
}, |
|
|
|
}, |
|
|
|
Template: core.PodTemplateSpec{ |
|
|
|
Template: core.PodTemplateSpec{ |
|
|
|
ObjectMeta: meta.ObjectMeta{ |
|
|
|
ObjectMeta: meta.ObjectMeta{ |
|
|
|
Labels: map[string]string{ |
|
|
|
Labels: labels.Set{ |
|
|
|
"app": name, |
|
|
|
"app": name, |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNameLabel: svc.Name, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
|
svcNamespaceLabel: svc.Namespace, |
|
|
@ -442,6 +496,25 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
Spec: core.PodSpec{ |
|
|
|
Spec: core.PodSpec{ |
|
|
|
ServiceAccountName: "svclb", |
|
|
|
ServiceAccountName: "svclb", |
|
|
|
AutomountServiceAccountToken: utilpointer.Bool(false), |
|
|
|
AutomountServiceAccountToken: utilpointer.Bool(false), |
|
|
|
|
|
|
|
SecurityContext: &core.PodSecurityContext{ |
|
|
|
|
|
|
|
Sysctls: sysctls, |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
Tolerations: []core.Toleration{ |
|
|
|
|
|
|
|
{ |
|
|
|
|
|
|
|
Key: "node-role.kubernetes.io/master", |
|
|
|
|
|
|
|
Operator: "Exists", |
|
|
|
|
|
|
|
Effect: "NoSchedule", |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
{ |
|
|
|
|
|
|
|
Key: "node-role.kubernetes.io/control-plane", |
|
|
|
|
|
|
|
Operator: "Exists", |
|
|
|
|
|
|
|
Effect: "NoSchedule", |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
{ |
|
|
|
|
|
|
|
Key: "CriticalAddonsOnly", |
|
|
|
|
|
|
|
Operator: "Exists", |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
}, |
|
|
|
}, |
|
|
|
}, |
|
|
|
}, |
|
|
|
}, |
|
|
|
UpdateStrategy: apps.DaemonSetUpdateStrategy{ |
|
|
|
UpdateStrategy: apps.DaemonSetUpdateStrategy{ |
|
|
@ -453,18 +526,6 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
}, |
|
|
|
}, |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
var sysctls []core.Sysctl |
|
|
|
|
|
|
|
for _, ipFamily := range svc.Spec.IPFamilies { |
|
|
|
|
|
|
|
switch ipFamily { |
|
|
|
|
|
|
|
case core.IPv4Protocol: |
|
|
|
|
|
|
|
sysctls = append(sysctls, core.Sysctl{Name: "net.ipv4.ip_forward", Value: "1"}) |
|
|
|
|
|
|
|
case core.IPv6Protocol: |
|
|
|
|
|
|
|
sysctls = append(sysctls, core.Sysctl{Name: "net.ipv6.conf.all.forwarding", Value: "1"}) |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
ds.Spec.Template.Spec.SecurityContext = &core.PodSecurityContext{Sysctls: sysctls} |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for _, port := range svc.Spec.Ports { |
|
|
|
for _, port := range svc.Spec.Ports { |
|
|
|
portName := fmt.Sprintf("lb-%s-%d", strings.ToLower(string(port.Protocol)), port.Port) |
|
|
|
portName := fmt.Sprintf("lb-%s-%d", strings.ToLower(string(port.Protocol)), port.Port) |
|
|
|
container := core.Container{ |
|
|
|
container := core.Container{ |
|
|
@ -492,14 +553,6 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
Name: "DEST_PROTO", |
|
|
|
Name: "DEST_PROTO", |
|
|
|
Value: string(port.Protocol), |
|
|
|
Value: string(port.Protocol), |
|
|
|
}, |
|
|
|
}, |
|
|
|
{ |
|
|
|
|
|
|
|
Name: "DEST_PORT", |
|
|
|
|
|
|
|
Value: strconv.Itoa(int(port.Port)), |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
{ |
|
|
|
|
|
|
|
Name: "DEST_IPS", |
|
|
|
|
|
|
|
Value: strings.Join(svc.Spec.ClusterIPs, " "), |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
}, |
|
|
|
}, |
|
|
|
SecurityContext: &core.SecurityContext{ |
|
|
|
SecurityContext: &core.SecurityContext{ |
|
|
|
Capabilities: &core.Capabilities{ |
|
|
|
Capabilities: &core.Capabilities{ |
|
|
@ -510,31 +563,36 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
}, |
|
|
|
}, |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
ds.Spec.Template.Spec.Containers = append(ds.Spec.Template.Spec.Containers, container) |
|
|
|
if localTraffic { |
|
|
|
} |
|
|
|
container.Env = append(container.Env, |
|
|
|
|
|
|
|
core.EnvVar{ |
|
|
|
// Add toleration to noderole.kubernetes.io/master=*:NoSchedule
|
|
|
|
Name: "DEST_PORT", |
|
|
|
masterToleration := core.Toleration{ |
|
|
|
Value: strconv.Itoa(int(port.NodePort)), |
|
|
|
Key: "node-role.kubernetes.io/master", |
|
|
|
}, |
|
|
|
Operator: "Exists", |
|
|
|
core.EnvVar{ |
|
|
|
Effect: "NoSchedule", |
|
|
|
Name: "DEST_IPS", |
|
|
|
} |
|
|
|
ValueFrom: &core.EnvVarSource{ |
|
|
|
ds.Spec.Template.Spec.Tolerations = append(ds.Spec.Template.Spec.Tolerations, masterToleration) |
|
|
|
FieldRef: &core.ObjectFieldSelector{ |
|
|
|
|
|
|
|
FieldPath: "status.hostIP", |
|
|
|
// Add toleration to noderole.kubernetes.io/control-plane=*:NoSchedule
|
|
|
|
}, |
|
|
|
controlPlaneToleration := core.Toleration{ |
|
|
|
}, |
|
|
|
Key: "node-role.kubernetes.io/control-plane", |
|
|
|
}, |
|
|
|
Operator: "Exists", |
|
|
|
) |
|
|
|
Effect: "NoSchedule", |
|
|
|
} else { |
|
|
|
} |
|
|
|
container.Env = append(container.Env, |
|
|
|
ds.Spec.Template.Spec.Tolerations = append(ds.Spec.Template.Spec.Tolerations, controlPlaneToleration) |
|
|
|
core.EnvVar{ |
|
|
|
|
|
|
|
Name: "DEST_PORT", |
|
|
|
|
|
|
|
Value: strconv.Itoa(int(port.Port)), |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
core.EnvVar{ |
|
|
|
|
|
|
|
Name: "DEST_IPS", |
|
|
|
|
|
|
|
Value: strings.Join(svc.Spec.ClusterIPs, " "), |
|
|
|
|
|
|
|
}, |
|
|
|
|
|
|
|
) |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
// Add toleration to CriticalAddonsOnly
|
|
|
|
ds.Spec.Template.Spec.Containers = append(ds.Spec.Template.Spec.Containers, container) |
|
|
|
criticalAddonsOnlyToleration := core.Toleration{ |
|
|
|
|
|
|
|
Key: "CriticalAddonsOnly", |
|
|
|
|
|
|
|
Operator: "Exists", |
|
|
|
|
|
|
|
} |
|
|
|
} |
|
|
|
ds.Spec.Template.Spec.Tolerations = append(ds.Spec.Template.Spec.Tolerations, criticalAddonsOnlyToleration) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Add node selector only if label "svccontroller.k3s.cattle.io/enablelb" exists on the nodes
|
|
|
|
// Add node selector only if label "svccontroller.k3s.cattle.io/enablelb" exists on the nodes
|
|
|
|
enableNodeSelector, err := k.nodeHasDaemonSetLabel() |
|
|
|
enableNodeSelector, err := k.nodeHasDaemonSetLabel() |
|
|
@ -551,6 +609,7 @@ func (k *k3s) newDaemonSet(svc *core.Service) (*apps.DaemonSet, error) { |
|
|
|
} |
|
|
|
} |
|
|
|
ds.Labels[nodeSelectorLabel] = "true" |
|
|
|
ds.Labels[nodeSelectorLabel] = "true" |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
return ds, nil |
|
|
|
return ds, nil |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
@ -563,7 +622,7 @@ func (k *k3s) updateDaemonSets() error { |
|
|
|
return err |
|
|
|
return err |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
|
|
|
|
nodeSelector := labels.SelectorFromSet(map[string]string{nodeSelectorLabel: fmt.Sprintf("%t", !enableNodeSelector)}) |
|
|
|
nodeSelector := labels.SelectorFromSet(labels.Set{nodeSelectorLabel: fmt.Sprintf("%t", !enableNodeSelector)}) |
|
|
|
daemonsets, err := k.daemonsetCache.List(k.LBNamespace, nodeSelector) |
|
|
|
daemonsets, err := k.daemonsetCache.List(k.LBNamespace, nodeSelector) |
|
|
|
if err != nil { |
|
|
|
if err != nil { |
|
|
|
return err |
|
|
|
return err |
|
|
|