diff --git a/charts/kks-provider-plugin/Chart.yaml b/charts/kks-provider-plugin/Chart.yaml index f6f720f..6277e37 100644 --- a/charts/kks-provider-plugin/Chart.yaml +++ b/charts/kks-provider-plugin/Chart.yaml @@ -2,8 +2,8 @@ apiVersion: v2 name: kks-provider-plugin description: Combined Kloud provider plugin chart (LoadBalancer + CSI) for kks clusters type: application -version: 1.2.4 -appVersion: "1.1.1" +version: 1.2.5 +appVersion: "1.1.2" kubeVersion: ">=1.28.0-0" home: https://github.com/KubelanCloud/kks-provider-plugin sources: diff --git a/pkg/kloudlb/constants.go b/pkg/kloudlb/constants.go index 5bf38c0..74d031d 100644 --- a/pkg/kloudlb/constants.go +++ b/pkg/kloudlb/constants.go @@ -3,6 +3,7 @@ package kloudlb const ( Finalizer = "lb.kloude.ir/finalizer" AnnotationIP = "lb.kloude.ir/ip" + AnnotationVIP = "lb.kloude.ir/vip" AnnotationLBID = "lb.kloude.ir/id" Namespace = "kube-system" diff --git a/pkg/kloudlb/controller/controller.go b/pkg/kloudlb/controller/controller.go index fe51697..be94215 100644 --- a/pkg/kloudlb/controller/controller.go +++ b/pkg/kloudlb/controller/controller.go @@ -3,6 +3,7 @@ package controller import ( "context" "fmt" + "net" "os" "strings" "time" @@ -14,9 +15,9 @@ import ( "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" corelisters "k8s.io/client-go/listers/core/v1" + "k8s.io/client-go/tools/cache" "k8s.io/client-go/tools/leaderelection" "k8s.io/client-go/tools/leaderelection/resourcelock" - "k8s.io/client-go/tools/cache" "k8s.io/client-go/util/workqueue" "github.com/KubelanCloud/kks-provider-plugin/pkg/kloudlb" @@ -177,16 +178,22 @@ func (c *Controller) sync(ctx context.Context, key string) error { if err != nil { return err } + vip := strings.TrimSpace(lb.IP) + ingressIP, err := vipIngressIP(vip) + if err != nil { + return err + } patch := svc.DeepCopy() if patch.Annotations == nil { patch.Annotations = map[string]string{} } - patch.Annotations[kloudlb.AnnotationIP] = lb.IP + patch.Annotations[kloudlb.AnnotationIP] = ingressIP + patch.Annotations[kloudlb.AnnotationVIP] = vip patch.Annotations[kloudlb.AnnotationLBID] = lb.ID patch.Status = corev1.ServiceStatus{ LoadBalancer: corev1.LoadBalancerStatus{ - Ingress: []corev1.LoadBalancerIngress{{IP: lb.IP}}, + Ingress: []corev1.LoadBalancerIngress{{IP: ingressIP}}, }, } @@ -194,6 +201,29 @@ func (c *Controller) sync(ctx context.Context, key string) error { return err } +func vipIngressIP(vip string) (string, error) { + vip = strings.TrimSpace(vip) + if vip == "" { + return "", fmt.Errorf("load balancer vip is empty") + } + if strings.Contains(vip, "/") { + ip, _, err := net.ParseCIDR(vip) + if err != nil { + return "", fmt.Errorf("invalid load balancer vip %q: %w", vip, err) + } + ipv4 := ip.To4() + if ipv4 == nil { + return "", fmt.Errorf("load balancer vip %q is not ipv4", vip) + } + return ipv4.String(), nil + } + ip := net.ParseIP(vip) + if ip == nil || ip.To4() == nil { + return "", fmt.Errorf("invalid load balancer vip %q", vip) + } + return ip.String(), nil +} + func (c *Controller) finalize(ctx context.Context, svc *corev1.Service) error { if !containsString(svc.Finalizers, kloudlb.Finalizer) { return nil diff --git a/pkg/kloudlb/speaker/speaker.go b/pkg/kloudlb/speaker/speaker.go index d94813d..ef15bdf 100644 --- a/pkg/kloudlb/speaker/speaker.go +++ b/pkg/kloudlb/speaker/speaker.go @@ -5,11 +5,12 @@ import ( "fmt" "net" "os" + "strings" "sync" "time" - corev1 "k8s.io/api/core/v1" coordinationv1 "k8s.io/api/coordination/v1" + corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/wait" @@ -141,26 +142,31 @@ func (s *Speaker) sync(ctx context.Context, key string) error { return err } - ip := serviceExternalIP(svc) - if ip == "" || svc.DeletionTimestamp != nil { - s.stopLeading(key, ip) + vip := serviceExternalVIP(svc) + if vip == "" || svc.DeletionTimestamp != nil { + s.stopLeading(key, vip) return nil } leaseName := leaseNameFor(key) - leader, err := s.acquireLease(ctx, leaseName, ip) + leader, err := s.acquireLease(ctx, leaseName, vip) if err != nil { return err } if leader != s.nodeName { - s.stopLeading(key, ip) + s.stopLeading(key, vip) return nil } - return s.ensureVIP(ip) + return s.ensureVIP(vip) } -func serviceExternalIP(svc *corev1.Service) string { +func serviceExternalVIP(svc *corev1.Service) string { + if svc.Annotations != nil { + if vip := strings.TrimSpace(svc.Annotations[kloudlb.AnnotationVIP]); vip != "" { + return vip + } + } for _, ing := range svc.Status.LoadBalancer.Ingress { if ing.IP != "" { return ing.IP @@ -287,17 +293,15 @@ func addVIP(ifaceName, ip string) error { if err != nil { return fmt.Errorf("lookup interface %s: %w", ifaceName, err) } - parsed := net.ParseIP(ip) - if parsed == nil { - return fmt.Errorf("invalid vip %q", ip) - } - addr := &netlink.Addr{ - IPNet: &net.IPNet{IP: parsed.To4(), Mask: net.CIDRMask(32, 32)}, + ipNet, err := vipToIPNet(link, ip) + if err != nil { + return err } + addr := &netlink.Addr{IPNet: ipNet} if err := netlink.AddrAdd(link, addr); err != nil { return fmt.Errorf("add vip %s on %s: %w", ip, ifaceName, err) } - return sendGratuitousARP(link, parsed.To4()) + return sendGratuitousARP(link, ipNet.IP.To4()) } func removeVIP(ifaceName, ip string) error { @@ -305,16 +309,63 @@ func removeVIP(ifaceName, ip string) error { if err != nil { return err } - parsed := net.ParseIP(ip) - if parsed == nil { - return fmt.Errorf("invalid vip %q", ip) - } - addr := &netlink.Addr{ - IPNet: &net.IPNet{IP: parsed.To4(), Mask: net.CIDRMask(32, 32)}, + ipNet, err := vipToIPNet(link, ip) + if err != nil { + return err } + addr := &netlink.Addr{IPNet: ipNet} return netlink.AddrDel(link, addr) } +func vipToIPNet(link netlink.Link, vip string) (*net.IPNet, error) { + vip = strings.TrimSpace(vip) + if vip == "" { + return nil, fmt.Errorf("vip is required") + } + if strings.Contains(vip, "/") { + ip, ipNet, err := net.ParseCIDR(vip) + if err != nil { + return nil, fmt.Errorf("invalid vip %q: %w", vip, err) + } + ipv4 := ip.To4() + if ipv4 == nil { + return nil, fmt.Errorf("vip %q is not ipv4", vip) + } + return &net.IPNet{IP: ipv4, Mask: ipNet.Mask}, nil + } + + ip := net.ParseIP(vip) + if ip == nil || ip.To4() == nil { + return nil, fmt.Errorf("invalid vip %q", vip) + } + mask, err := interfaceIPv4Mask(link) + if err != nil { + return nil, err + } + return &net.IPNet{IP: ip.To4(), Mask: mask}, nil +} + +func interfaceIPv4Mask(link netlink.Link) (net.IPMask, error) { + addrs, err := netlink.AddrList(link, 0) + if err != nil { + return nil, fmt.Errorf("list ipv4 addresses on %s: %w", link.Attrs().Name, err) + } + for _, addr := range addrs { + if addr.IPNet == nil || addr.IPNet.IP == nil || addr.IPNet.Mask == nil { + continue + } + ip := addr.IPNet.IP.To4() + if ip == nil { + continue + } + if !ip.IsGlobalUnicast() { + continue + } + return addr.IPNet.Mask, nil + } + return nil, fmt.Errorf("cannot infer ipv4 mask from interface %s", link.Attrs().Name) +} + func sendGratuitousARP(link netlink.Link, ip net.IP) error { if ip == nil { return fmt.Errorf("ip is required")