Added common lease for multiple services for all modes and service election for BGP

Signed-off-by: Patryk Strusiewicz-Surmacki <patryk.pawel.strusiewicz-surmacki@external.telekom.de>
This commit is contained in:
Patryk Strusiewicz-Surmacki
2025-08-13 11:27:41 +02:00
committed by Marcel Fest
parent b1183e8a93
commit 95995500bc
19 changed files with 1540 additions and 232 deletions

View File

@@ -39,7 +39,7 @@ jobs:
if: matrix.mode== 'bgp'
- name: Change log directory permissions
run: sudo chmod -R 755 /tmp/kube-vip-test-${{ matrix.mode }}*
if: matrix.mode== 'bgp'
if: matrix.mode== 'bgp' && always()
- name: Save logs
uses: actions/upload-artifact@v4
with:
@@ -67,5 +67,5 @@ jobs:
uses: actions/upload-artifact@v4
with:
name: services-test-logs-${{ steps.date.outputs.date }}
path: /tmp/kube-vip-services*
path: /tmp/kube-vip-service-tests*
if: always()

View File

@@ -2,6 +2,7 @@ package arp
import (
"context"
"fmt"
log "log/slog"
"sync"
"time"
@@ -11,7 +12,7 @@ import (
)
type Manager struct {
instances map[string]*Instance
instances sync.Map
config *kubevip.Config
}
@@ -24,8 +25,7 @@ type Instance struct {
func NewManager(config *kubevip.Config) *Manager {
return &Manager{
instances: make(map[string]*Instance),
config: config,
config: config,
}
}
@@ -42,10 +42,14 @@ func (i *Instance) Name() string {
}
func (m *Manager) Insert(instance *Instance) {
i, ok := m.instances[instance.Name()]
if !ok {
log.Info("inserting ARP/NDP instance", "name", instance.Name())
m.instances[instance.Name()] = instance
i, err := m.get(instance.Name())
if err != nil {
log.Error("[ARP manager] unable to insert instance", "err", err)
return
}
if i == nil {
log.Info("[ARP manager] inserting ARP/NDP instance", "name", instance.Name())
m.instances.Store(instance.Name(), instance)
} else {
i.mu.Lock()
defer i.mu.Unlock()
@@ -54,20 +58,35 @@ func (m *Manager) Insert(instance *Instance) {
}
func (m *Manager) Remove(instance *Instance) {
if i, ok := m.instances[instance.Name()]; ok {
i, err := m.get(instance.Name())
if err != nil {
log.Error("[ARP manager] unable to remove the instance", "err", err)
return
}
if i != nil {
i.mu.Lock()
defer i.mu.Unlock()
if i.counter > 1 {
i.counter--
} else {
log.Info("removing ARP/NDP instance", "name", instance.Name())
delete(m.instances, instance.Name())
log.Info("[ARP manager] removing ARP/NDP instance", "name", instance.Name())
if _, err := instance.network.DeleteIP(); err != nil {
log.Error("failed to delete IP", "address", instance.network.IP(), "err", err)
}
m.instances.Delete(instance.Name())
}
} else {
log.Warn("[ARP manager] unable to remove the instance - instance not found", "name", instance.Name())
}
}
func (m *Manager) Count(name string) int {
if i, ok := m.instances[name]; ok {
i, err := m.get(name)
if err != nil {
log.Error("[ARP manager] unable to count instance", "err", err)
return -1
}
if i != nil {
i.mu.Lock()
defer i.mu.Unlock()
return i.counter
@@ -76,26 +95,46 @@ func (m *Manager) Count(name string) int {
}
func (m *Manager) StartAdvertisement(ctx context.Context) {
log.Info("Starting ARP/NDP advertisement")
log.Info("[ARP manager] starting ARP/NDP advertisement")
for {
select {
case <-ctx.Done(): // if cancel() execute
return
default:
for _, instance := range m.instances {
if instance.counter > 0 {
ensureIPAndSendGratuitous(instance)
m.instances.Range(func(_ any, instance any) bool {
if i, ok := instance.(*Instance); ok {
if i.counter > 0 {
ensureIPAndSendGratuitous(i)
} else {
// this instance should not be advertised - delete the IP just in case...
if _, err := i.network.DeleteIP(); err != nil {
log.Error("[ARP manager] failed to delete IP", "address", i.network.IP(), "err", err)
}
}
}
}
return true
})
}
if m.config.ArpBroadcastRate < 500 {
log.Error("arp broadcast rate is too low", "rate (ms)", m.config.ArpBroadcastRate, "setting to (ms)", "3000")
log.Warn("[ARP manager] arp broadcast rate is too low", "rate (ms)", m.config.ArpBroadcastRate, "setting to (ms)", "3000")
m.config.ArpBroadcastRate = 3000
}
time.Sleep(time.Duration(m.config.ArpBroadcastRate) * time.Millisecond)
}
}
func (m *Manager) get(name string) (*Instance, error) {
i, exists := m.instances.Load(name)
if !exists {
return nil, nil
}
inst, ok := i.(*Instance)
if !ok {
return nil, fmt.Errorf("value for name %q is not of Instance pointer type", name)
}
return inst, nil
}
// ensureIPAndSendGratuitous - adds IP to the interface if missing, and send
// either a gratuitous ARP or gratuitous NDP. Re-adds the interface if it is IPv6
// and in a dadfailed state.
@@ -112,20 +151,14 @@ func ensureIPAndSendGratuitous(instance *Instance) {
}
if deleted {
log.Info("deleted and recreating address", "IP", ipString, "interface", iface)
// if _, err := instance.network.AddIP(false); err != nil {
// log.Error("failed to recreate address", "IP", ipString, "interface", iface)
// }
}
}
// Ensure the address exists on the interface before attempting to ARP
// if instance.network.HasEndpoints() {
if added, err := instance.network.AddIP(true); err != nil {
log.Warn(err.Error())
} else if added {
log.Warn("Re-applied the VIP configuration", "ip", ipString, "interface", iface)
}
// }
if vip.IsIPv6(ipString) {
// Gratuitous NDP, will broadcast new MAC <-> IPv6 address
@@ -145,5 +178,4 @@ func ensureIPAndSendGratuitous(instance *Instance) {
log.Warn(err.Error())
}
}
}

View File

@@ -340,7 +340,7 @@ func (cluster *Cluster) StartLoadBalancerService(ctx context.Context, c *kubevip
return
}
for i := range cluster.Network {
if c.EnableARP && cluster.arpMgr.Count(cluster.Network[i].ARPName()) > 0 {
if c.EnableARP && cluster.arpMgr.Count(cluster.Network[i].ARPName()) > 1 {
continue
}
log.Info("[VIP] Deleting VIP", "ip", cluster.Network[i].IP())

View File

@@ -110,6 +110,10 @@ func ClearBGPHosts(service *v1.Service, instances *[]*instance.Instance, bgpServ
}
func ClearBGPHostsByInstance(instance *instance.Instance, bgpServer *bgp.Server) {
if instance == nil {
log.Error("failed to clear BGP host for nil instance")
return
}
for _, cluster := range instance.Clusters {
for i := range cluster.Network {
network := cluster.Network[i]

View File

@@ -117,9 +117,9 @@ func (rt *RoutingTable) deleteAction(service *v1.Service) {
}
func (rt *RoutingTable) setInstanceEndpointsStatus(service *v1.Service, endpoints []string) error {
instance := instance.FindServiceInstance(service, *rt.instances)
instance := instance.FindServiceInstanceWithTimeout(service, *rt.instances)
if instance == nil {
log.Error("failed to find the instance", "service", service.UID, "provider", rt.provider.GetLabel())
log.Error("failed to find the instance", "namespace", service.Namespace, "name", service.Name, "uid", service.UID, "provider", rt.provider.GetLabel())
} else {
for _, c := range instance.Clusters {
for n := range c.Network {
@@ -140,23 +140,34 @@ func (rt *RoutingTable) setInstanceEndpointsStatus(service *v1.Service, endpoint
func ClearRoutes(service *v1.Service, instances *[]*instance.Instance) []error {
errs := []error{}
if instance := instance.FindServiceInstance(service, *instances); instance != nil {
for _, cluster := range instance.Clusters {
for i := range cluster.Network {
route := cluster.Network[i].PrepareRoute()
// check if route we are about to delete is not referenced by more than one service
if CountRouteReferences(route, instances) <= 1 {
err := cluster.Network[i].DeleteRoute()
if err != nil && !errors.Is(err, syscall.ESRCH) {
log.Error("failed to delete route", "ip", cluster.Network[i].IP(), "err", err)
errs = append(errs, err)
}
log.Debug("deleted route", "ip",
cluster.Network[i].IP(), "service name", service.Name, "namespace", service.Namespace, "interface", cluster.Network[i].Interface())
if svcInst := instance.FindServiceInstance(service, *instances); svcInst != nil {
clearErrs := ClearRoutesByInstance(service, svcInst, instances)
errs = append(errs, clearErrs...)
}
return errs
}
func ClearRoutesByInstance(service *v1.Service, svcInst *instance.Instance, instances *[]*instance.Instance) []error {
if svcInst == nil {
return []error{fmt.Errorf("failed to remove routes for nil instance of service %s/%s, uid: %s", service.Namespace, service.Name, service.UID)}
}
errs := []error{}
for _, cluster := range svcInst.Clusters {
for i := range cluster.Network {
route := cluster.Network[i].PrepareRoute()
// check if route we are about to delete is not referenced by more than one service
if CountRouteReferences(route, instances) <= 1 {
err := cluster.Network[i].DeleteRoute()
if err != nil && !errors.Is(err, syscall.ESRCH) {
log.Error("failed to delete route", "ip", cluster.Network[i].IP(), "err", err)
errs = append(errs, err)
}
log.Debug("deleted route", "ip",
cluster.Network[i].IP(), "service name", service.Name, "namespace", service.Namespace, "interface", cluster.Network[i].Interface())
}
}
}
return errs
}

View File

@@ -5,6 +5,7 @@ import (
"net"
"strconv"
"strings"
"time"
"log/slog"
log "log/slog"
@@ -478,5 +479,27 @@ func FindServiceInstance(svc *v1.Service, instances []*Instance) *Instance {
return instances[i]
}
}
log.Debug("insance not found", "UID", svc.UID)
return nil
}
func FindServiceInstanceWithTimeout(svc *v1.Service, instances []*Instance) *Instance {
log.Debug("finding service with timeout", "namespace", svc.Namespace, "name", svc.Name, "UID", svc.UID)
ticker := time.NewTicker(time.Millisecond * 200)
defer ticker.Stop()
to := time.NewTimer(time.Second * 60)
defer to.Stop()
for {
select {
case <-to.C:
return nil
case <-ticker.C:
for i := range instances {
log.Debug("saved service", "instance", i, "UID", instances[i].ServiceSnapshot.UID)
if instances[i].ServiceSnapshot.UID == svc.UID {
return instances[i]
}
}
}
}
}

View File

@@ -46,4 +46,6 @@ const (
UpnpEnabled = "kube-vip.io/forwardUPNP"
RPFilter = "kube-vip.io/rp_filter" // Set the return path filter for a specific service interface
ServiceLease = "kube-vip.io/leaseName"
)

127
pkg/lease/lease.go Normal file
View File

@@ -0,0 +1,127 @@
package lease
import (
"context"
"fmt"
"sync"
"github.com/kube-vip/kube-vip/pkg/kubevip"
v1 "k8s.io/api/core/v1"
)
// Manager is used to manage leases.
type Manager struct {
leases map[string]*Lease
lock sync.Mutex
}
// NewManager creates new lease manager.
func NewManager() *Manager {
return &Manager{
leases: make(map[string]*Lease),
}
}
// Add adds lease or incerements counter if lease is alreay used.
func (m *Manager) Add(service *v1.Service) (*Lease, bool) {
m.lock.Lock()
defer m.lock.Unlock()
_, id := GetName(service)
if _, exist := m.leases[id]; !exist {
ctx, cancel := context.WithCancel(context.Background())
m.leases[id] = newLease(ctx, cancel)
return m.leases[id], true
}
m.leases[id].increment()
return m.leases[id], false
}
// Delete decrements lease counter and removes the lease if counter equals 0.
func (m *Manager) Delete(service *v1.Service) {
m.lock.Lock()
defer m.lock.Unlock()
_, id := GetName(service)
if _, exist := m.leases[id]; exist {
m.leases[id].decrement()
if m.leases[id].cnt < 1 {
delete(m.leases, id)
}
}
}
// Get returns lease for the service.
func (m *Manager) Get(service *v1.Service) *Lease {
m.lock.Lock()
defer m.lock.Unlock()
_, id := GetName(service)
if lease, exist := m.leases[id]; exist {
return lease
}
return nil
}
// GetLeaderContext returns leder context for the service.
func (m *Manager) GetLeaderContext(service *v1.Service) context.Context {
m.lock.Lock()
defer m.lock.Unlock()
_, id := GetName(service)
if _, ok := m.leases[id]; !ok {
return nil
}
return m.leases[id].Ctx
}
// Lease holds lease data.
type Lease struct {
cnt uint
Lock *sync.Mutex
Ctx context.Context
Cancel context.CancelFunc
Started chan any
}
func newLease(ctx context.Context, cancel context.CancelFunc) *Lease {
return &Lease{
Ctx: ctx,
Cancel: cancel,
cnt: 1,
Lock: new(sync.Mutex),
Started: make(chan any),
}
}
func (l *Lease) increment() {
l.Lock.Lock()
defer l.Lock.Unlock()
l.cnt++
}
func (l *Lease) decrement() {
l.Lock.Lock()
defer l.Lock.Unlock()
if l.cnt == 0 {
return
}
l.cnt--
if l.cnt < 1 {
l.Cancel()
}
}
// GetName gets lease name and id for the service.
func GetName(service *v1.Service) (string, string) {
serviceLease, exists := service.Annotations[kubevip.ServiceLease]
if !exists || serviceLease == "" {
serviceLease = fmt.Sprintf("kubevip-%s", service.Name)
}
serviceLeaseID := fmt.Sprintf("%s/%s", serviceLease, service.Namespace)
return serviceLease, serviceLeaseID
}
// UsesCommon checks if service uses common lease feature.
func UsesCommon(service *v1.Service) bool {
_, common := service.Annotations[kubevip.ServiceLease]
return common
}

View File

@@ -105,9 +105,18 @@ func (sm *Manager) startBGP() error {
}
}
err = sm.svcProcessor.ServicesWatcher(ctx, sm.svcProcessor.SyncServices)
if err != nil {
return err
if sm.config.EnableServicesElection {
log.Info("beginning watching services, leaderelection will happen for every service")
err = sm.svcProcessor.StartServicesWatchForLeaderElection(ctx)
if err != nil {
return err
}
} else {
log.Info("beginning watching services without leader election")
err = sm.svcProcessor.ServicesWatcher(ctx, sm.svcProcessor.SyncServices)
if err != nil {
return err
}
}
log.Info("Shutting down Kube-Vip")

View File

@@ -1,36 +0,0 @@
package services
import (
"context"
"sync"
)
type Context struct {
Ctx context.Context
Cancel context.CancelFunc
IsActive bool
IsWatched bool
ConfiguredNetworks sync.Map
}
func NewContext(ctx context.Context) *Context {
svcCtx, svcCancel := context.WithCancel(ctx)
return &Context{
Ctx: svcCtx,
Cancel: svcCancel,
}
}
func (ctx *Context) HasConfiguredNetworks() bool {
cnt := 0
ctx.ConfiguredNetworks.Range(func(_ any, _ any) bool {
cnt++
return cnt < 1
})
return cnt > 0
}
func (ctx *Context) IsNetworkConfigured(ip string) bool {
_, exists := ctx.ConfiguredNetworks.Load(ip)
return exists
}

View File

@@ -3,25 +3,17 @@ package services
import (
"context"
"fmt"
"sync"
"time"
log "log/slog"
"github.com/kube-vip/kube-vip/pkg/lease"
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/tools/leaderelection"
"k8s.io/client-go/tools/leaderelection/resourcelock"
)
var (
svcLocks map[string]*sync.Mutex
)
func init() {
svcLocks = make(map[string]*sync.Mutex)
}
// The StartServicesWatchForLeaderElection function will start a services watcher, the
func (p *Processor) StartServicesWatchForLeaderElection(ctx context.Context) error {
err := p.ServicesWatcher(ctx, p.StartServicesLeaderElection)
@@ -29,12 +21,14 @@ func (p *Processor) StartServicesWatchForLeaderElection(ctx context.Context) err
return err
}
for _, instance := range p.ServiceInstances {
for _, cluster := range instance.Clusters {
for i := range cluster.Network {
_ = cluster.Network[i].DeleteRoute()
if p.config.EnableRoutingTable {
for _, instance := range p.ServiceInstances {
for _, cluster := range instance.Clusters {
for i := range cluster.Network {
_ = cluster.Network[i].DeleteRoute()
}
cluster.Stop()
}
cluster.Stop()
}
}
@@ -45,8 +39,7 @@ func (p *Processor) StartServicesWatchForLeaderElection(ctx context.Context) err
// The startServicesWatchForLeaderElection function will start a services watcher, the
func (p *Processor) StartServicesLeaderElection(ctx context.Context, service *v1.Service) error {
serviceLease := fmt.Sprintf("kubevip-%s", service.Name)
serviceLeaseID := fmt.Sprintf("%s/%s", serviceLease, service.Namespace)
serviceLease, _ := lease.GetName(service)
log.Info("new leader election", "service", service.Name, "namespace", service.Namespace, "lock_name", serviceLease, "host_id", p.config.NodeName)
// we use the Lease lock type since edits to Leases are less common
// and fewer objects in the cluster watch "all Leases".
@@ -60,15 +53,12 @@ func (p *Processor) StartServicesLeaderElection(ctx context.Context, service *v1
Identity: p.config.NodeName,
},
}
childCtx, childCancel := context.WithCancel(ctx)
defer childCancel()
if _, ok := svcLocks[serviceLeaseID]; !ok {
svcLocks[serviceLeaseID] = new(sync.Mutex)
}
svcLocks[serviceLeaseID].Lock()
defer svcLocks[serviceLeaseID].Unlock()
go func() {
// wait for the service context to end and delete the lease then
<-ctx.Done()
p.leaseMgr.Delete(service)
}()
svcCtx, err := p.getServiceContext(service.UID)
if err != nil {
@@ -80,8 +70,40 @@ func (p *Processor) StartServicesLeaderElection(ctx context.Context, service *v1
svcCtx.IsActive = true
svcLease, isNew := p.leaseMgr.Add(service)
// this service is sharing lease
if !isNew {
// wait for leader election to start or context to be done
select {
case <-svcLease.Started:
case <-svcLease.Ctx.Done():
svcCtx.IsActive = false
return nil
}
<-svcLease.Started
if lease.UsesCommon(service) {
if err := p.SyncServices(ctx, service); err != nil {
log.Error("service sync", "err", err)
svcLease.Cancel()
}
// just block until context is cancelled
<-ctx.Done()
if svcCtx.IsActive {
if err := p.deleteService(service.UID); err != nil {
log.Error("service deletion", "err", err)
}
}
}
// wait for leaderelection to be finished
<-svcLease.Ctx.Done()
// Mark this service is inactive
svcCtx.IsActive = false
return nil
}
// start the leader election code loop
leaderelection.RunOrDie(childCtx, leaderelection.LeaderElectionConfig{
leaderelection.RunOrDie(svcLease.Ctx, leaderelection.LeaderElectionConfig{
Lock: lock,
// IMPORTANT: you MUST ensure that any code you have that
// is protected by the lease must terminate **before**
@@ -99,8 +121,9 @@ func (p *Processor) StartServicesLeaderElection(ctx context.Context, service *v1
// we run this in background as it's blocking
if err := p.SyncServices(ctx, service); err != nil {
log.Error("service sync", "err", err)
childCancel()
svcLease.Cancel()
}
close(svcLease.Started)
},
OnStoppedLeading: func() {
// we can do cleanup here

View File

@@ -13,6 +13,7 @@ import (
"github.com/kube-vip/kube-vip/pkg/endpoints/providers"
"github.com/kube-vip/kube-vip/pkg/instance"
"github.com/kube-vip/kube-vip/pkg/kubevip"
"github.com/kube-vip/kube-vip/pkg/lease"
"github.com/kube-vip/kube-vip/pkg/networkinterface"
"github.com/kube-vip/kube-vip/pkg/servicecontext"
"github.com/kube-vip/kube-vip/pkg/vip"
@@ -46,6 +47,8 @@ type Processor struct {
intfMgr *networkinterface.Manager
arpMgr *arp.Manager
leaseMgr *lease.Manager
}
func NewServicesProcessor(config *kubevip.Config, bgpServer *bgp.Server,
@@ -71,8 +74,9 @@ func NewServicesProcessor(config *kubevip.Config, bgpServer *bgp.Server,
Help: "Count all events fired by the service watcher categorised by event type",
}, []string{"type"}),
intfMgr: intfMgr,
arpMgr: arpMgr,
intfMgr: intfMgr,
arpMgr: arpMgr,
leaseMgr: lease.NewManager(),
}
}
@@ -106,6 +110,12 @@ func (p *Processor) AddOrModify(ctx context.Context, event watch.Event, serviceF
return true, nil
}
_, usesCommonLease := svc.Annotations[kubevip.ServiceLease]
if usesCommonLease && svc.Spec.ExternalTrafficPolicy != v1.ServiceExternalTrafficPolicyTypeCluster {
return false, fmt.Errorf("annotation %q cannot be used with service traffic policy other than %q",
kubevip.ServiceLease, v1.ServiceExternalTrafficPolicyTypeCluster)
}
svcCtx, err := p.getServiceContext(svc.UID)
if err != nil {
return false, fmt.Errorf("failed to get service context: %w", err)
@@ -138,7 +148,6 @@ func (p *Processor) AddOrModify(ctx context.Context, event watch.Event, serviceF
originalService := instance.FetchServiceAddresses(i.ServiceSnapshot)
newService := instance.FetchServiceAddresses(svc)
if !reflect.DeepEqual(originalService, newService) {
// Calls the cancel function of the context
if svcCtx != nil {
log.Warn("(svcs) The load balancer has changed, cancelling original load balancer")
@@ -147,8 +156,7 @@ func (p *Processor) AddOrModify(ctx context.Context, event watch.Event, serviceF
<-svcCtx.Ctx.Done()
}
err = p.deleteService(svc.UID)
if err != nil {
if err := p.deleteService(svc.UID); err != nil {
log.Error("(svc) unable to remove", "service", svc.UID)
}
@@ -266,13 +274,14 @@ func (p *Processor) AddOrModify(ctx context.Context, event watch.Event, serviceF
func (p *Processor) Delete(event watch.Event) (bool, error) {
svc, ok := event.Object.(*v1.Service)
if !ok {
return false, fmt.Errorf("unable to parse Kubernetes services from API watcher")
return false, fmt.Errorf("(svcs) unable to parse Kubernetes services from API watcher")
}
svcCtx, err := p.getServiceContext(svc.UID)
if err != nil {
return false, fmt.Errorf("(svcs) unable to get context: %w", err)
}
if svcCtx != nil && svcCtx.IsActive {
if svcCtx != nil {
// We only care about LoadBalancer services
if svc.Spec.Type != v1.ServiceTypeLoadBalancer {
return true, nil
@@ -280,7 +289,7 @@ func (p *Processor) Delete(event watch.Event) (bool, error) {
// We can ignore this service
if svc.Annotations["kube-vip.io/ignore"] == "true" {
log.Info("(svcs)ignore annotation for kube-vip", "service name", svc.Name)
log.Info("(svcs) ignore annotation for kube-vip", "service name", svc.Name)
return true, nil
}
@@ -306,14 +315,6 @@ func (p *Processor) Delete(event watch.Event) (bool, error) {
p.svcMap.Delete(svc.UID)
}
if p.config.EnableLeaderElection && !p.config.EnableServicesElection {
if p.config.EnableBGP {
endpoints.ClearBGPHosts(svc, &p.ServiceInstances, p.bgpServer)
} else if p.config.EnableRoutingTable {
endpoints.ClearRoutes(svc, &p.ServiceInstances)
}
}
log.Info("(svcs) deleted", "service name", svc.Name, "namespace", svc.Namespace)
return true, nil

View File

@@ -12,6 +12,7 @@ import (
"github.com/google/go-cmp/cmp"
"github.com/vishvananda/netlink"
v1 "k8s.io/api/core/v1"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/wait"
@@ -42,19 +43,19 @@ func (p *Processor) SyncServices(ctx context.Context, svc *v1.Service) error {
action := p.getServiceInstanceAction(svc)
switch action {
case ActionDelete:
log.Debug("[service] delete", "namespace", svc.Namespace, "name", svc.Name)
log.Debug("[service] delete", "namespace", svc.Namespace, "name", svc.Name, "uid", svc.UID)
if err := p.deleteService(svc.UID); err != nil {
return fmt.Errorf("error deleting service %s/%s: %w", svc.Namespace, svc.Name, err)
}
case ActionAdd:
log.Debug("[service] add", "namespace", svc.Namespace, "name", svc.Name)
log.Debug("[service] add", "namespace", svc.Namespace, "name", svc.Name, "uid", svc.UID)
if err := p.addService(ctx, svc); err != nil {
return fmt.Errorf("error adding service %s/%s: %w", svc.Namespace, svc.Name, err)
}
case ActionNone:
log.Debug("[service] no action", "namespace", svc.Namespace, "name", svc.Name)
log.Debug("[service] no action", "namespace", svc.Namespace, "name", svc.Name, "uid", svc.UID)
}
log.Debug("[FINISHED] Service Sync", "namespace", svc.Namespace, "name", svc.Name)
log.Debug("[FINISHED] Service Sync", "namespace", svc.Namespace, "name", svc.Name, "uid", svc.UID)
return nil
}
@@ -85,7 +86,7 @@ func (p *Processor) getServiceInstanceAction(svc *v1.Service) ServiceInstanceAct
return ActionDelete
}
}
if !comparePortsAndPortStatuses(svc) {
if len(svc.Status.LoadBalancer.Ingress) > 0 && !comparePortsAndPortStatuses(svc) {
return ActionDelete
}
}
@@ -94,7 +95,7 @@ func (p *Processor) getServiceInstanceAction(svc *v1.Service) ServiceInstanceAct
}
}
if len(addresses) > 0 {
log.Debug("No matching service instance found", "service", svc.Name, "namespace", svc.Namespace, "addresses", addresses)
log.Debug("no matching service instance found", "service", svc.Name, "namespace", svc.Namespace, "addresses", addresses)
return ActionAdd // If no matching instance is found, we need to add a new service instance
}
return ActionNone
@@ -266,10 +267,12 @@ func (p *Processor) deleteService(uid types.UID) error {
serviceInstance = p.ServiceInstances[x]
}
}
// If we've been through all services and not found the correct one then error
if !found {
// TODO: - fix UX
// return fmt.Errorf("unable to find/stop service [%s]", uid)
log.Error("unable to find/stop service", "uid", uid)
return nil
}
@@ -287,6 +290,7 @@ func (p *Processor) deleteService(uid types.UID) error {
vipSet[vip] = nil
}
}
for _, vip := range instance.FetchServiceAddresses(serviceInstance.ServiceSnapshot) {
if _, found := vipSet[vip]; found {
shared = true
@@ -308,9 +312,15 @@ func (p *Processor) deleteService(uid types.UID) error {
return fmt.Errorf("[service] error deleting DHCP Link : %v", err)
}
}
for i := range serviceInstance.VIPConfigs {
if serviceInstance.VIPConfigs[i].EnableBGP {
endpoints.ClearBGPHostsByInstance(serviceInstance, p.bgpServer)
if p.config.EnableBGP {
endpoints.ClearBGPHostsByInstance(serviceInstance, p.bgpServer)
}
if p.config.EnableRoutingTable && (p.config.EnableLeaderElection || p.config.EnableServicesElection) {
if errs := endpoints.ClearRoutesByInstance(serviceInstance.ServiceSnapshot, serviceInstance, &p.ServiceInstances); len(errs) > 0 {
for _, err := range errs {
log.Error("unable to clear routes", "err", err)
}
}
}
@@ -329,7 +339,7 @@ func (p *Processor) deleteService(uid types.UID) error {
// Update the service array
p.ServiceInstances = updatedInstances
log.Info("Removed instance from manager", "uid", uid, "remaining advertised services", len(p.ServiceInstances))
log.Info("Removed instance from manager", "uid", uid, "name", serviceInstance.ServiceSnapshot.Name, "remaining advertised services", len(p.ServiceInstances))
return nil
}
@@ -484,7 +494,7 @@ func (p *Processor) updateStatus(i *instance.Instance) error {
if !cmp.Equal(currentService.Status.LoadBalancer.Ingress, ingresses) {
currentService.Status.LoadBalancer.Ingress = ingresses
_, err = p.clientSet.CoreV1().Services(currentService.Namespace).UpdateStatus(context.TODO(), currentService, metav1.UpdateOptions{})
if err != nil {
if err != nil && !apierrors.IsInvalid(err) {
log.Error("updating Service", "namespace", i.ServiceSnapshot.Namespace, "name", i.ServiceSnapshot.Name, "err", err)
return err
}

13
test.sh
View File

@@ -1,13 +0,0 @@
#!/usr/bin/bash
rm -f logs.txt
for i in $(seq 1 100);
do
echo "RUN $i"
GOMAXPROCS=4 make e2e-tests129
if [[ "$?" -ne 0 ]]; then
echo "FAILED AT RUN $i"
break
fi
done

View File

@@ -147,10 +147,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath,
k8sImagePath, e2e.IPv4Family, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129,
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4")
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
@@ -166,7 +167,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv4 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, "")
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -174,7 +176,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, "")
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -200,10 +203,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath,
k8sImagePath, e2e.IPv6Family, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129,
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4")
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
@@ -217,7 +221,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, "")
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -225,7 +230,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, "")
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -251,10 +257,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamily, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "fixed", "mpbgp-ipv4")
logger, nodesNumber, "fixed", "mpbgp-ipv4", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
@@ -268,7 +275,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv6 routes over IPv4 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv6)
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv6, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -276,7 +284,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv6)
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv6, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -302,10 +311,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamilyIPv6, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "fixed", "mpbgp-ipv6")
logger, nodesNumber, "fixed", "mpbgp-ipv6", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, n := range gobgpPeers {
@@ -319,7 +329,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv4 routes over IPv6 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv4)
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv4, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -327,7 +338,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv4)
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv4, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -354,10 +366,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
_, containerIP = setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamily, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv4")
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv4", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
@@ -371,7 +384,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv6 routes over IPv4 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, containerIP)
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -379,7 +393,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, containerIP)
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -406,10 +421,11 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
Expect(err).NotTo(HaveOccurred())
containerIP, _ = setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamilyIPv6, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv6")
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv6", false)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, n := range gobgpPeers {
@@ -423,7 +439,8 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("advertise IPv4 routes over IPv6 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 1, gobgpClient, containerIP)
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -431,29 +448,407 @@ var _ = Describe("kube-vip BGP mode", Ordered, func() {
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName, trafficPolicy, client, 2, gobgpClient, containerIP)
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
})
Describe("kube-vip IPv4 services BGP mode functionality", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath,
k8sImagePath, e2e.IPv4Family, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129,
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
Eventually(func() error {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: p.IP,
})
return err
}, "30s", "200ms").Should(Succeed())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv4 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip IPv6 services BGP mode functionality", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath,
k8sImagePath, e2e.IPv6Family, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129,
kubeVIPBGPManifestTemplate, &gobgpClient, logger, nodesNumber, "", "bgp-ipv4", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: p.IP,
})
Expect(err).ToNot(HaveOccurred())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, "", "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services BGP mode functionality with MP-BGP IPv6 over IPv4 - fixed nexthop", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamily, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "fixed", "mpbgp-ipv4", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: p.IP,
})
Expect(err).ToNot(HaveOccurred())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv6 routes over IPv4 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv6, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv6, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were delete while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv6, "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services BGP mode functionality with MP-BGP IPv4 over IPv6 - fixed nexthop", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamilyIPv6, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "fixed", "mpbgp-ipv6", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, n := range gobgpPeers {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: n.IP,
})
Expect(err).ToNot(HaveOccurred())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv4 routes over IPv6 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, defaultFixedNexthopv4, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv4, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, defaultFixedNexthopv4, "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services BGP mode functionality with MP-BGP IPv6 over IPv4 - auto_sourceif nexthop", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
containerIP string
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
_, containerIP = setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamily, e2e.IPv4Family, []string{e2e.IPv4Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv4", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, p := range gobgpPeers {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: p.IP,
})
Expect(err).ToNot(HaveOccurred())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv6 routes over IPv4 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv6Family, api.Family_AFI_IP6, []corev1.IPFamily{corev1.IPv6Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services BGP mode functionality with MP-BGP IPv4 over IPv6 - fixed nexthop", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
gobgpClient api.GobgpApiClient
gobgpPeers []*e2e.BGPPeerValues
containerIP string
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
containerIP, _ = setupEnv(&cpVIP, &clusterName, tempDirPath, manifestValues, localIPv4, localIPv6, imagePath, configPath, k8sImagePath,
e2e.DualstackFamilyIPv6, e2e.IPv6Family, []string{e2e.IPv6Family}, &client, &gobgpPeers, v129, kubeVIPBGPManifestTemplate, &gobgpClient,
logger, nodesNumber, "auto_sourceif", "mpbgp-if-ipv6", true)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
for _, n := range gobgpPeers {
_, err := gobgpClient.DeletePeer(context.TODO(), &api.DeletePeerRequest{
Address: n.IP,
})
Expect(err).ToNot(HaveOccurred())
}
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("advertise IPv4 routes over IPv6 session",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 1, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only stops advertising route if it was referenced by multiple services and all of them were deleted while using common lease",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
testBGP(offset, e2e.IPv4Family, api.Family_AFI_IP, []corev1.IPFamily{corev1.IPv4Protocol}, svcName,
trafficPolicy, client, 2, gobgpClient, containerIP, "common-lease")
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
}
})
func testBGP(offset uint, lbFamily string, afiFamily api.Family_Afi, svcFamily []corev1.IPFamily, svcName string,
trafficPolicy corev1.ServiceExternalTrafficPolicy, client kubernetes.Interface, numberOfServices int, gobgpClient api.GobgpApiClient, expectedNexthop string) {
trafficPolicy corev1.ServiceExternalTrafficPolicy, client kubernetes.Interface, numberOfServices int,
gobgpClient api.GobgpApiClient, expectedNexthop string, serviceLease string) {
lbAddress := e2e.GenerateVIP(lbFamily, offset)
routeCheckFamily := &api.Family{
Afi: afiFamily,
Safi: api.Family_SAFI_UNICAST,
}
testServiceBGP(svcName, lbAddress, trafficPolicy, client, svcFamily, numberOfServices, gobgpClient, routeCheckFamily, expectedNexthop)
testServiceBGP(svcName, lbAddress, trafficPolicy, client, svcFamily, numberOfServices, gobgpClient, routeCheckFamily, expectedNexthop, serviceLease)
}
func setupEnv(cpVIP, clusterName *string, tempDirPath string, manifestValues *e2e.KubevipManifestValues,
localIPv4, localIPv6, imagePath, configPath, k8sImagePath, clusterAddrFamily, bgpClientAddrFamily string, peerAddrFamily []string, client *kubernetes.Interface,
gobgpPeers *[]*e2e.BGPPeerValues, v129 bool, kubeVIPBGPManifestTemplate *template.Template, gobgpClient *api.GobgpApiClient,
logger log.Logger, nodesNumber int, mpbgpnexthop, clusterNameSuffix string) (string, string) {
logger log.Logger, nodesNumber int, mpbgpnexthop, clusterNameSuffix string, serviceElection bool) (string, string) {
var err error
*cpVIP = e2e.GenerateVIP(clusterAddrFamily, SOffset.Get())
@@ -510,7 +905,7 @@ func setupEnv(cpVIP, clusterName *string, tempDirPath string, manifestValues *e2
ConfigPath: configPath,
ControlPlaneEnable: "false",
SvcEnable: "true",
SvcElectionEnable: "false",
SvcElectionEnable: fmt.Sprintf("%t", serviceElection),
BGPAS: kubevipAS,
BGPPeers: strings.Join(kvPeersStr, ","),
MPBGPNexthop: mpbgpnexthop,
@@ -520,7 +915,8 @@ func setupEnv(cpVIP, clusterName *string, tempDirPath string, manifestValues *e2
By(manifestValues.BGPPeers)
*clusterName, *client, _ = prepareCluster(tempDirPath, clusterNameSuffix, k8sImagePath, v129, kubeVIPBGPManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
*clusterName, *client, _ = prepareCluster(tempDirPath, clusterNameSuffix, k8sImagePath, v129,
kubeVIPBGPManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
container := fmt.Sprintf("%s-control-plane", *clusterName)
@@ -592,7 +988,8 @@ func setupEnv(cpVIP, clusterName *string, tempDirPath string, manifestValues *e2
func testServiceBGP(svcName, lbAddress string, trafficPolicy corev1.ServiceExternalTrafficPolicy,
client kubernetes.Interface, serviceAddrFamily []corev1.IPFamily, numberOfServices int,
gobgpClient api.GobgpApiClient, gobgpFamily *api.Family, expectedNexthop string) {
gobgpClient api.GobgpApiClient, gobgpFamily *api.Family, expectedNexthop string,
serviceLease string) {
lbAddresses := vip.Split(lbAddress)
services := []string{}
@@ -602,11 +999,15 @@ func testServiceBGP(svcName, lbAddress string, trafficPolicy corev1.ServiceExter
for _, svc := range services {
createTestService(svc, dsNamespace, dsName, lbAddress,
client, corev1.IPFamilyPolicyPreferDualStack, serviceAddrFamily, trafficPolicy)
client, corev1.IPFamilyPolicyPreferDualStack, serviceAddrFamily, trafficPolicy, serviceLease, 80)
time.Sleep(time.Second)
}
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking bgp route for address %q", addr)))
paths := checkGoBGPPaths(context.Background(), gobgpClient, gobgpFamily, []*api.TableLookupPrefix{{Prefix: addr}}, 1)
Expect(paths).ToNot(BeNil())
Expect(paths).ToNot(BeEmpty())
Expect(strings.Contains(paths[0].Prefix, lbAddress)).To(BeTrue())
if expectedNexthop != "" {
Expect(strings.Contains(paths[0].String(), fmt.Sprintf("next_hop:\"%s\"", expectedNexthop)) || strings.Contains(paths[0].String(), fmt.Sprintf("next_hops:\"%s\"", expectedNexthop))).To(BeTrue())
@@ -614,11 +1015,16 @@ func testServiceBGP(svcName, lbAddress string, trafficPolicy corev1.ServiceExter
}
for i := range numberOfServices {
By(withTimestamp(fmt.Sprintf("deleting service '%s/%s'", dsNamespace, services[i])))
err := client.CoreV1().Services(dsNamespace).Delete(context.TODO(), services[i], metav1.DeleteOptions{})
Expect(err).ToNot(HaveOccurred())
time.Sleep(time.Second)
if i < numberOfServices-1 {
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking bgp route for address %q", addr)))
paths := checkGoBGPPaths(context.Background(), gobgpClient, gobgpFamily, []*api.TableLookupPrefix{{Prefix: addr}}, 1)
Expect(paths).ToNot(BeNil())
Expect(paths).ToNot(BeEmpty())
Expect(strings.Contains(paths[0].Prefix, lbAddress)).To(BeTrue())
if expectedNexthop != "" {
Expect(strings.Contains(paths[0].String(), fmt.Sprintf("next_hop:\"%s\"", expectedNexthop)) || strings.Contains(paths[0].String(), fmt.Sprintf("next_hops:\"%s\"", expectedNexthop))).To(BeTrue())
@@ -628,6 +1034,7 @@ func testServiceBGP(svcName, lbAddress string, trafficPolicy corev1.ServiceExter
}
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking bgp route for address %q - should be deleted", addr)))
checkGoBGPPaths(context.Background(), gobgpClient, gobgpFamily, []*api.TableLookupPrefix{{Prefix: addr}}, 0)
}
}
@@ -668,6 +1075,25 @@ func newGoBGPClient(address string, port uint32) (api.GobgpApiClient, error) {
}
func checkGoBGPPaths(ctx context.Context, client api.GobgpApiClient, family *api.Family, prefixes []*api.TableLookupPrefix, expectedPaths int) []*api.Destination {
// ticker := time.NewTicker(time.Second)
// defer ticker.Stop()
// to := time.NewTimer(time.Second * 180)
// defer to.Stop()
// for {
// select {
// case <-to.C:
// return nil
// case <-ticker.C:
// paths, err := getGoBGPPaths(ctx, client, family, prefixes)
// if err != nil {
// return nil
// }
// if len(paths) == expectedPaths {
// return paths
// }
// }
// }
var paths []*api.Destination
Eventually(func() error {
var err error
@@ -679,7 +1105,7 @@ func checkGoBGPPaths(ctx context.Context, client api.GobgpApiClient, family *api
return fmt.Errorf("expected %d paths, but found %d", expectedPaths, len(paths))
}
return nil
}, "120s").ShouldNot(HaveOccurred())
}, "360s", "1s").ShouldNot(HaveOccurred())
return paths
}

View File

@@ -101,10 +101,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ipv4", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -159,10 +161,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ipv6", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -228,10 +232,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-ipv4", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, addSAN)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, addSAN, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -303,10 +309,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-ipv6", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, addSAN)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, addSAN, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -372,10 +380,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv4", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -384,7 +394,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("configures an IPv4 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv4Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName, trafficPolicy, client, svcElection, ipFamily, 1)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -393,7 +404,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv4Family, offset)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName, trafficPolicy, client, svcElection, ipFamily, 2)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -438,10 +450,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv6", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -450,7 +464,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("configures an IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName, trafficPolicy, client, svcElection, ipFamily, 1)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -459,7 +474,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName, trafficPolicy, client, svcElection, ipFamily, 2)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -505,10 +521,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv4", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -517,7 +535,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("configures an IPv4 and IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName, trafficPolicy, client, svcElection, ipFamily, 1)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -526,7 +545,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName, trafficPolicy, client, svcElection, ipFamily, 2)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -574,10 +594,12 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv6", k8sImagePath, v129, kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
@@ -586,7 +608,8 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("configures an IPv4 and IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName, trafficPolicy, client, svcElection, ipFamily, 1)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
@@ -595,17 +618,338 @@ var _ = Describe("kube-vip routing table mode", Ordered, func() {
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName, trafficPolicy, client, svcElection, ipFamily, 2)
testServiceRT(svcName, lbAddress, "plndr-svcs-lock", "kube-system", clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
})
Describe("kube-vip IPv4 services routing table mode functionality with services election", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
svcElection bool
ipFamily []corev1.IPFamily
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
cpVIP = e2e.GenerateVIP(e2e.IPv4Family, SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.IPv4Family,
}
manifestValues = &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
ControlPlaneEnable: "false",
SvcEnable: "true",
SvcElectionEnable: "true",
}
var err error
svcElection, err = strconv.ParseBool(manifestValues.SvcElectionEnable)
Expect(err).ToNot(HaveOccurred())
ipFamily = []corev1.IPFamily{corev1.IPv4Protocol}
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an IPv4 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv4Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv4Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted if common lease is used",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv4Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, true)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip IPv6 services routing table mode functionality with services election", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
svcElection bool
ipFamily []corev1.IPFamily
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
cpVIP = e2e.GenerateVIP(e2e.IPv6Family, SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.IPv6Family,
}
manifestValues = &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
ControlPlaneEnable: "false",
SvcEnable: "true",
SvcElectionEnable: "true",
}
var err error
svcElection, err = strconv.ParseBool(manifestValues.SvcElectionEnable)
Expect(err).ToNot(HaveOccurred())
ipFamily = []corev1.IPFamily{corev1.IPv6Protocol}
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-svc-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted when common lease is used",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, true)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services routing table mode functionality - IPv4 primary with services election", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
svcElection bool
ipFamily []corev1.IPFamily
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
cpVIP = e2e.GenerateDualStackVIP(SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.DualStackFamily,
}
manifestValues = &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
ControlPlaneEnable: "false",
SvcEnable: "true",
SvcElectionEnable: "true",
EnableEndpointslices: "true",
}
var err error
svcElection, err = strconv.ParseBool(manifestValues.SvcElectionEnable)
Expect(err).ToNot(HaveOccurred())
ipFamily = []corev1.IPFamily{corev1.IPv4Protocol, corev1.IPv6Protocol}
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv4", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an IPv4 and IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted when common lease is used",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, true)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack services routing table mode functionality - IPv6 primary with services election", Ordered, func() {
var (
cpVIP string
clusterName string
client kubernetes.Interface
manifestValues *e2e.KubevipManifestValues
svcElection bool
ipFamily []corev1.IPFamily
tempDirPath string
nodesNumber = 1
)
BeforeAll(func() {
cpVIP = e2e.GenerateDualStackVIP(SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.DualStackFamily,
PodSubnet: "fd00:10:244::/56,10.244.0.0/16",
ServiceSubnet: "fd00:10:96::/112,10.96.0.0/16",
}
manifestValues = &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
ControlPlaneEnable: "false",
SvcEnable: "true",
SvcElectionEnable: "true",
EnableEndpointslices: "true",
}
var err error
svcElection, err = strconv.ParseBool(manifestValues.SvcElectionEnable)
Expect(err).ToNot(HaveOccurred())
ipFamily = []corev1.IPFamily{corev1.IPv4Protocol, corev1.IPv6Protocol}
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "rt-ds-svc-ipv6", k8sImagePath, v129,
kubeVIPRoutingTableManifestTemplate, logger, manifestValues, networking, nodesNumber, nil, 1)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
err := e2e.GetLogs(context.Background(), client, tempDirPath)
Expect(err).ToNot(HaveOccurred())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an IPv4 and IPv6 routes for services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 1, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, false)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
Entry("with external traffic policy - local", "test-svc-local", SOffset.Get(), corev1.ServiceExternalTrafficPolicyLocal),
)
DescribeTable("only removes route if it was referenced by multiple services and all of them were deleted when common lease is used",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceRT(svcName, lbAddress, fmt.Sprintf("kubevip-%s", svcName), dsNamespace, clusterName,
trafficPolicy, client, svcElection, ipFamily, 2, true)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
}
})
func testServiceRT(svcName, lbAddress, leaseName, leaseNamespace, clusterName string, trafficPolicy corev1.ServiceExternalTrafficPolicy,
client kubernetes.Interface, serviceElection bool, ipFamily []corev1.IPFamily, numberOfServices int) {
client kubernetes.Interface, serviceElection bool, ipFamily []corev1.IPFamily, numberOfServices int, commonLease bool) {
lbAddresses := vip.Split(lbAddress)
services := []string{}
@@ -614,28 +958,51 @@ func testServiceRT(svcName, lbAddress, leaseName, leaseNamespace, clusterName st
}
for _, svc := range services {
svcLease := ""
if commonLease {
svcLease = leaseName
}
createTestService(svc, dsNamespace, dsName, lbAddress,
client, corev1.IPFamilyPolicyPreferDualStack, ipFamily, trafficPolicy)
client, corev1.IPFamilyPolicyPreferDualStack, ipFamily, trafficPolicy, svcLease, 1)
time.Sleep(time.Second)
}
var container string
if serviceElection {
container = e2e.GetLeaseHolder(leaseName, leaseNamespace, client)
for _, svc := range services {
lease := fmt.Sprintf("kubevip-%s", svc)
if commonLease {
lease = leaseName
}
By(withTimestamp(fmt.Sprintf("getting lease holder for lease '%s/%s'", leaseNamespace, lease)))
container = e2e.GetLeaseHolder(lease, leaseNamespace, client)
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking route presence for address %q on container %q", addr, container)))
e2e.CheckRoutePresence(addr, container, true)
}
}
} else {
container = fmt.Sprintf("%s-control-plane", clusterName)
}
for _, addr := range lbAddresses {
e2e.CheckRoutePresence(addr, container, true)
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking route presence for address %q on container %q", addr, container)))
e2e.CheckRoutePresence(addr, container, true)
}
}
for i := range numberOfServices {
By(withTimestamp(fmt.Sprintf("deleting service %s/%s\n", dsNamespace, services[i])))
expected := i < numberOfServices-1
err := client.CoreV1().Services(dsNamespace).Delete(context.TODO(), services[i], metav1.DeleteOptions{})
Expect(err).ToNot(HaveOccurred())
time.Sleep(time.Second)
for _, addr := range lbAddresses {
By(withTimestamp(fmt.Sprintf("checking route presence for address %q on container %q - expected: %t", addr, container, expected)))
e2e.CheckRoutePresence(addr, container, expected)
if commonLease {
By(withTimestamp(fmt.Sprintf("getting lease holder for lease '%s/%s' - expected: %t", leaseNamespace, leaseName, expected)))
e2e.CheckLeasePresence(leaseName, leaseNamespace, client, expected)
}
}
}
}

View File

@@ -12,6 +12,7 @@ import (
"os"
"os/exec"
"path/filepath"
"strconv"
"strings"
"sync"
"text/template"
@@ -127,7 +128,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil, 1)
})
AfterAll(func() {
@@ -170,7 +172,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "svc-ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -227,7 +230,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -276,7 +280,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, clusterCfg = prepareCluster(tempDirPath, "ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil)
clusterName, client, clusterCfg = prepareCluster(tempDirPath, "ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil, 1)
})
AfterAll(func() {
@@ -321,7 +326,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "svc-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -378,7 +384,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -426,7 +433,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil, 1)
})
AfterAll(func() {
@@ -469,7 +477,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -527,7 +536,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-el-ipv4", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-el-ipv4", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -578,7 +588,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 3, nil, 1)
})
AfterAll(func() {
@@ -624,7 +635,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -690,7 +702,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-el-ipv6", k8sImagePath, v129, kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ds-svc-el-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 1, nil, 1)
})
AfterAll(func() {
@@ -739,7 +752,8 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "ipv4-hostname", k8sImagePath, v129, kubeVIPHostnameManifestTemplate, logger, manifestValues, networking, 3, nil)
clusterName, client, _ = prepareCluster(tempDirPath, "ipv4-hostname", k8sImagePath, v129,
kubeVIPHostnameManifestTemplate, logger, manifestValues, networking, 3, nil, 1)
})
AfterAll(func() {
@@ -754,6 +768,110 @@ var _ = Describe("kube-vip ARP/NDP broadcast neighbor", Ordered, func() {
testControlPlaneVIPs([]string{cpVIP}, clusterName, client)
})
})
Describe("kube-vip IPv6 functionality with manually specified lease name, svc_enable=true, svc_election=true", Ordered, func() {
var (
cpVIP string
client kubernetes.Interface
clusterName string
tempDirPath string
)
BeforeAll(func() {
cpVIP = e2e.GenerateVIP(e2e.IPv6Family, SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.IPv6Family,
}
manifestValues := &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
SvcEnable: "true",
SvcElectionEnable: "true",
EnableEndpointslices: "false",
}
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-m-ipv6", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 2, nil, 3)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
Eventually(func() error {
return e2e.GetLogs(context.Background(), client, tempDirPath)
}, "60s", "5s").Should(Succeed())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an single IPv6 VIP address for multiple services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateVIP(e2e.IPv6Family, offset)
testServiceCommonLease(svcName, lbAddress, dsNamespace,
trafficPolicy,
client, []corev1.IPFamily{corev1.IPv6Protocol}, 3)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
Describe("kube-vip DualStack functionality with manually specified lease name, svc_enable=true, svc_election=true", Ordered, func() {
var (
cpVIP string
client kubernetes.Interface
clusterName string
tempDirPath string
)
BeforeAll(func() {
cpVIP = e2e.GenerateDualStackVIP(SOffset.Get())
networking := &kindconfigv1alpha4.Networking{
IPFamily: kindconfigv1alpha4.DualStackFamily,
PodSubnet: "fd00:10:244::/56,10.244.0.0/16",
ServiceSubnet: "fd00:10:96::/112,10.96.0.0/16",
}
manifestValues := &e2e.KubevipManifestValues{
ControlPlaneVIP: cpVIP,
ImagePath: imagePath,
ConfigPath: configPath,
SvcEnable: "true",
SvcElectionEnable: "true",
EnableEndpointslices: "false",
}
var err error
tempDirPath, err = os.MkdirTemp(tempDirPathRoot, "kube-vip-test")
Expect(err).NotTo(HaveOccurred())
clusterName, client, _ = prepareCluster(tempDirPath, "svc-el-m-ds", k8sImagePath, v129,
kubeVIPManifestTemplate, logger, manifestValues, networking, 2, nil, 3)
})
AfterAll(func() {
By(fmt.Sprintf("saving logs to %q", tempDirPath))
Eventually(func() error {
return e2e.GetLogs(context.Background(), client, tempDirPath)
}, "60s", "5s").Should(Succeed())
cleanupCluster(clusterName, ConfigMtx, logger)
})
DescribeTable("configures an single IPv4 and IPv6 VIP address for multiple services",
func(svcName string, offset uint, trafficPolicy corev1.ServiceExternalTrafficPolicy) {
lbAddress := e2e.GenerateDualStackVIP(offset)
testServiceCommonLease(svcName, lbAddress, dsNamespace,
trafficPolicy,
client, []corev1.IPFamily{corev1.IPv6Protocol, corev1.IPv4Protocol}, 3)
},
Entry("with external traffic policy - cluster", "test-svc-cluster", SOffset.Get(), corev1.ServiceExternalTrafficPolicyCluster),
)
})
}
})
@@ -812,13 +930,13 @@ func assertConnection(protocol, ip, port, suffix string, transportTimeout, event
}
defer resp.Body.Close()
return resp.StatusCode
}, eventuallyTimeout).Should(BeElementOf([]int{http.StatusOK, http.StatusInternalServerError}), fmt.Sprintf("Failed to connect to %s", ip))
}, eventuallyTimeout, transportTimeout).Should(BeElementOf([]int{http.StatusOK, http.StatusInternalServerError}), fmt.Sprintf("Failed to connect to %s", ip))
By(withTimestamp(fmt.Sprintf("estabilished connection to %s://%s:%s/%s", protocol, ip, port, suffix)))
}
// Assume connection to the provided address is possible
func assertConnectionError(protocol, ip, port, suffix string, transportTimeout time.Duration) {
func assertConnectionError(protocol, ip, port, suffix string, transportTimeout time.Duration, eventuallyTimeout time.Duration) {
if strings.Contains(ip, ":") {
ip = fmt.Sprintf("[%s]", ip)
}
@@ -833,12 +951,12 @@ func assertConnectionError(protocol, ip, port, suffix string, transportTimeout t
Eventually(func() error {
_, err := client.Get(fmt.Sprintf("%s://%s:%s/%s", protocol, ip, port, suffix))
return err
}, time.Second*30).Should(HaveOccurred())
}, eventuallyTimeout, transportTimeout).Should(HaveOccurred())
By(withTimestamp(fmt.Sprintf("connection %s://%s:%s/%s error as expected", protocol, ip, port, suffix)))
}
func killLeader(leaderName string, clusterName string) {
func killLeader(leaderName string) {
cmd := exec.Command(
"docker", "kill", leaderName,
)
@@ -875,7 +993,7 @@ func withTimestamp(text string) string {
return fmt.Sprintf("%s: %s", time.Now(), text)
}
func createTestDS(name, namespace string, client kubernetes.Interface) {
func createTestDS(name, namespace string, client kubernetes.Interface, port int) {
labels := make(map[string]string)
labels["app"] = name
d := v1.DaemonSet{
@@ -911,13 +1029,13 @@ func createTestDS(name, namespace string, client kubernetes.Interface) {
Image: "ghcr.io/traefik/whoami:v1.11",
Ports: []corev1.ContainerPort{
{
ContainerPort: 80,
ContainerPort: int32(port),
},
},
Env: []corev1.EnvVar{
{
Name: "PORT",
Value: "80",
Name: "WHOAMI_PORT_NUMBER",
Value: strconv.Itoa(port),
},
},
},
@@ -931,9 +1049,13 @@ func createTestDS(name, namespace string, client kubernetes.Interface) {
Expect(err).ToNot(HaveOccurred())
}
func createTestService(name, namespace, target, lbAddress string, client kubernetes.Interface, ipfPolicy corev1.IPFamilyPolicy, ipFamiles []corev1.IPFamily, externalPolicy corev1.ServiceExternalTrafficPolicy) {
func createTestService(name, namespace, target, lbAddress string, client kubernetes.Interface, ipfPolicy corev1.IPFamilyPolicy,
ipFamiles []corev1.IPFamily, externalPolicy corev1.ServiceExternalTrafficPolicy, leaseName string, port int) {
svcAnnotations := make(map[string]string)
svcAnnotations[kubevip.LoadbalancerIPAnnotation] = lbAddress
if leaseName != "" {
svcAnnotations[kubevip.ServiceLease] = leaseName
}
labels := make(map[string]string)
labels["app"] = target
@@ -955,7 +1077,7 @@ func createTestService(name, namespace, target, lbAddress string, client kuberne
Ports: []corev1.ServicePort{
{
Protocol: corev1.ProtocolTCP,
Port: 80,
Port: int32(port),
},
},
Selector: labels,
@@ -967,9 +1089,6 @@ func createTestService(name, namespace, target, lbAddress string, client kuberne
return err
}, time.Second*60, time.Second).Should(Succeed())
By(withTimestamp(fmt.Sprintf("service %s/%s created", namespace, name)))
svcs, err := client.CoreV1().Services(namespace).List(context.Background(), metav1.ListOptions{})
Expect(err).ToNot(HaveOccurred())
By(svcs.String())
Eventually(func() error {
By(withTimestamp(fmt.Sprintf("getting service %s/%s\n", namespace, name)))
@@ -1000,7 +1119,7 @@ func checkIPAddressByLease(name, namespace, lbAddress string, expected bool, cli
By(withTimestamp(fmt.Sprintf("checking LB %q by lease %s/%s, should exist: %t", lbAddress, namespace, name, expected)))
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
to := time.NewTimer(time.Second * 180)
to := time.NewTimer(time.Second * 360)
defer to.Stop()
for {
select {
@@ -1017,7 +1136,7 @@ func checkIPAddressByLease(name, namespace, lbAddress string, expected bool, cli
func prepareCluster(tempDirPath, clusterNameSuffix, k8sImagePath string,
v129 bool, kubeVIPManifestTemplate *template.Template, logger log.Logger,
manifestValues *e2e.KubevipManifestValues, networking *kindconfigv1alpha4.Networking, nodesNum int,
addSAN *san) (string, kubernetes.Interface, *rest.Config) {
addSAN *san, dsNumber int) (string, kubernetes.Interface, *rest.Config) {
manifestPath := filepath.Join(tempDirPath, fmt.Sprintf("kube-vip-%s.yaml", clusterNameSuffix))
@@ -1091,7 +1210,13 @@ func prepareCluster(tempDirPath, clusterNameSuffix, k8sImagePath string,
client, cfg := createKindCluster(logger, &clusterConfig, clusterName)
By(withTimestamp("creating test daemonset"))
createTestDS(dsName, dsNamespace, client)
for i := range dsNumber {
tmpDsName := dsName
if i > 0 {
tmpDsName = fmt.Sprintf("%s-%d", dsName, i)
}
createTestDS(tmpDsName, dsNamespace, client, 80+i)
}
By(withTimestamp("loading local docker image to kind cluster"))
e2e.LoadDockerImageToKind(logger, manifestValues.ImagePath, clusterName)
@@ -1138,7 +1263,7 @@ func testControlPlaneVIPs(cpVIPs []string, clusterName string, client kubernetes
Eventually(client.CoreV1().Nodes().Delete(context.Background(), leaderName, metav1.DeleteOptions{}), "60s", "1s").Should(Succeed())
By(withTimestamp("killing the leader Kubernetes control plane node to trigger a fail-over scenario"))
killLeader(leaderName, clusterName)
killLeader(leaderName)
By(withTimestamp("checking that the Kubernetes control plane nodes are still accessible via the assigned VIP with little downtime"))
for _, cpVIP := range cpVIPs {
@@ -1166,7 +1291,8 @@ func testService(svcName, lbAddress, leaseName, leaseNamespace string, trafficPo
for _, svc := range services {
createTestService(svc, dsNamespace, dsName, lbAddress,
client, corev1.IPFamilyPolicyPreferDualStack, ipFamily, trafficPolicy)
client, corev1.IPFamilyPolicyPreferDualStack, ipFamily, trafficPolicy, "", 80)
time.Sleep(time.Second)
}
for _, addr := range lbAddresses {
@@ -1181,6 +1307,7 @@ func testService(svcName, lbAddress, leaseName, leaseNamespace string, trafficPo
err := client.CoreV1().Services(dsNamespace).Delete(context.TODO(), services[i], metav1.DeleteOptions{})
Expect(err).ToNot(HaveOccurred())
time.Sleep(time.Second)
for _, addr := range lbAddresses {
if serviceElection {
@@ -1192,7 +1319,71 @@ func testService(svcName, lbAddress, leaseName, leaseNamespace string, trafficPo
if expected {
assertConnection("http", addr, "80", "", 3*time.Second, 60*time.Second)
} else {
assertConnectionError("http", addr, "80", "", 3*time.Second)
assertConnectionError("http", addr, "80", "", 3*time.Second, 30*time.Second)
}
}
}
}
func testServiceCommonLease(svcName, lbAddress, leaseNamespace string, trafficPolicy corev1.ServiceExternalTrafficPolicy,
client kubernetes.Interface, ipFamily []corev1.IPFamily, numberOfServices int) {
lbAddresses := vip.Split(lbAddress)
services := []string{}
for i := range numberOfServices {
services = append(services, fmt.Sprintf("%s-%d", svcName, i))
}
lease := "common-lease"
for i, svc := range services {
tmpDsName := dsName
if i > 0 {
tmpDsName = fmt.Sprintf("%s-%d", tmpDsName, i)
}
createTestService(svc, dsNamespace, tmpDsName, lbAddress,
client, corev1.IPFamilyPolicyPreferDualStack, ipFamily, trafficPolicy, lease, 80+i)
time.Sleep(time.Second)
}
for i, addr := range lbAddresses {
checkIPAddressByLease(lease, leaseNamespace, addr, true, client)
assertConnection("http", addr, strconv.Itoa(80+i), "", 3*time.Second, 60*time.Second)
}
container := e2e.GetLeaseHolder(lease, leaseNamespace, client)
nodes, err := client.CoreV1().Nodes().List(context.Background(), metav1.ListOptions{})
Expect(err).ToNot(HaveOccurred())
for _, node := range nodes.Items {
expected := node.Name == container
for _, addr := range lbAddresses {
checkIPAddress(addr, node.Name, expected)
}
}
for i := range numberOfServices {
expected := i < numberOfServices-1
By(fmt.Sprintf("deleting service %q", services[i]))
err := client.CoreV1().Services(dsNamespace).Delete(context.TODO(), services[i], metav1.DeleteOptions{})
Expect(err).ToNot(HaveOccurred())
time.Sleep(time.Second)
Eventually(func() error {
_, err := client.CoreV1().Services(dsNamespace).Get(context.Background(), services[i], metav1.GetOptions{})
return err
}).ShouldNot(Succeed())
for _, addr := range lbAddresses {
for _, node := range nodes.Items {
if node.Name == container {
checkIPAddress(addr, node.Name, expected)
} else {
checkIPAddress(addr, node.Name, false)
}
}
}
}

View File

@@ -165,7 +165,7 @@ func CheckIPAddressPresence(ip string, container string, expected bool) bool {
func CheckIPAddressPresenceByLease(name, namespace, ip string, client kubernetes.Interface, expected bool) bool {
container := GetLeaseHolder(name, namespace, client)
By("Lease: " + container)
By("Lease holder: " + container)
if container == "" {
return false
}
@@ -191,23 +191,48 @@ func CheckRoutePresence(ip string, container string, expected bool) bool {
cmd.Stderr = cmdErr
cmd.Run()
By("Routes: " + cmdOut.String())
result = strings.Contains(cmdOut.String(), ip) == expected
return result
}, "60s", "1s").Should(BeTrue())
}, "120s", "1s").Should(BeTrue())
return result
}
func GetLeaseHolder(name, namespace string, client kubernetes.Interface) string {
func CheckLeasePresence(name, namespace string, client kubernetes.Interface, expected bool) *coordinationv1.Lease {
var lease *coordinationv1.Lease
Eventually(func() error {
var err error
lease, err = client.CoordinationV1().Leases(namespace).Get(context.TODO(), name, metav1.GetOptions{})
return err
}, "120s", "1s").ShouldNot(HaveOccurred())
}, "360s", "1s").ShouldNot(HaveOccurred())
Expect(lease).ToNot(BeNil())
if expected {
Expect(lease.Spec.HolderIdentity).ToNot(BeNil())
} else {
if lease.Spec.HolderIdentity != nil {
if expected {
Expect(*lease.Spec.HolderIdentity).ToNot(BeEmpty())
} else {
Expect(*lease.Spec.HolderIdentity).To(BeEmpty())
}
}
}
return lease
}
func GetLeaseHolder(name, namespace string, client kubernetes.Interface) string {
var lease *coordinationv1.Lease
Eventually(func() string {
lease = CheckLeasePresence(name, namespace, client, true)
return *lease.Spec.HolderIdentity
}, "360s", "1s").ShouldNot(BeEmpty())
return *lease.Spec.HolderIdentity
}

View File

@@ -7,10 +7,12 @@ import (
"fmt"
"net"
"net/http"
"os"
"strings"
"time"
"github.com/gookit/slog"
"github.com/kube-vip/kube-vip/testing/e2e"
"github.com/vishvananda/netlink"
"k8s.io/client-go/kubernetes"
@@ -29,12 +31,35 @@ import (
func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kubernetes.Clientset) []error {
var err error
var errs []error
globalTempDirPath, err := os.MkdirTemp("/tmp", "kube-vip-service-tests")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
defer func() {
if os.Getenv("E2E_KEEP_LOGS") != "true" {
if err := os.RemoveAll(globalTempDirPath); err != nil {
slog.Error(fmt.Errorf("failed to remove temporary directory %q: %w", globalTempDirPath, err))
}
}
}()
if config.Simple {
err = config.SimpleDeployment(ctx, clientset)
if err != nil {
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "Simple")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.Deployments {
@@ -43,6 +68,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "Deployments")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.LeaderFailover {
@@ -52,6 +86,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "LeaderFailover")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.LeaderActive {
@@ -61,6 +104,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "LeaderActive")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.LocalDeploy {
@@ -70,6 +122,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "LocalDeploy")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.Egress {
@@ -79,6 +140,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "EgressDeployment")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.Egress && config.EgressInternal {
@@ -88,6 +158,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "EgressDeploymentInternal")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.EgressIPv6 {
@@ -97,6 +176,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "EgressIPv6")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.EgressIPv6 && config.EgressInternal {
@@ -106,6 +194,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "EgressIPv6Internal")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
if config.DualStack {
@@ -115,6 +212,15 @@ func (config *TestConfig) StartServiceTest(ctx context.Context, clientset *kuber
slog.Error(err)
errs = append(errs, err)
}
tempDirPath, err := os.MkdirTemp(globalTempDirPath, "DualStack")
if err != nil {
slog.Error(err)
return []error{fmt.Errorf("failed to create temporary directory: %w", err)}
}
slog.Infof("saving logs to %q", tempDirPath)
if err := e2e.GetLogs(ctx, clientset, tempDirPath); err != nil {
slog.Error(err)
}
}
const testComplete = "🏆 Testing Complete [%d] passed / [%d] failed"