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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,7 @@ type Config struct {
KubernetesPathMode kubernetes.PathMode `yaml:"-"`
KubernetesNamespace string `yaml:"kubernetes-namespace"`
KubernetesEnableEndpointSlices bool `yaml:"enable-kubernetes-endpointslices"`
KubernetesTopologyZone string `yaml:"kubernetes-topology-zone"`
KubernetesEnableEastWest bool `yaml:"enable-kubernetes-east-west"`
KubernetesEastWestDomain string `yaml:"kubernetes-east-west-domain"`
KubernetesEastWestRangeDomains *listFlag `yaml:"kubernetes-east-west-range-domains"`
Expand Down Expand Up @@ -541,6 +542,7 @@ func NewConfig() *Config {
flag.StringVar(&cfg.KubernetesPathModeString, "kubernetes-path-mode", "kubernetes-ingress", "controls the default interpretation of Kubernetes ingress paths: <kubernetes-ingress|path-regexp|path-prefix>")
flag.StringVar(&cfg.KubernetesNamespace, "kubernetes-namespace", "", "watch only this namespace for ingresses")
flag.BoolVar(&cfg.KubernetesEnableEndpointSlices, "enable-kubernetes-endpointslices", false, "Enables that skipper fetches Kubernetes endpointslices instead of endpoints to scale more than 1000 pods within a service")
flag.StringVar(&cfg.KubernetesTopologyZone, "kubernetes-topology-zone", "", "sets the topology zone to be used for zone aware routing")
flag.BoolVar(&cfg.KubernetesEnableEastWest, "enable-kubernetes-east-west", false, "*Deprecated*: use kubernetes-east-west-range feature. Enables east-west communication, which automatically adds routes for Ingress objects with hostname <name>.<namespace>.skipper.cluster.local")
flag.StringVar(&cfg.KubernetesEastWestDomain, "kubernetes-east-west-domain", "", "*Deprecated*: use kubernetes-east-west-range feature. Sets the east-west domain, defaults to .skipper.cluster.local")
flag.Var(cfg.KubernetesEastWestRangeDomains, "kubernetes-east-west-range-domains", "set the cluster internal domains for east west traffic. Identified routes to such domains will include the -kubernetes-east-west-range-predicates")
Expand Down Expand Up @@ -988,6 +990,7 @@ func (c *Config) ToOptions() skipper.Options {
KubernetesPathMode: c.KubernetesPathMode,
KubernetesNamespace: c.KubernetesNamespace,
KubernetesEnableEndpointslices: c.KubernetesEnableEndpointSlices,
KubernetesTopologyZone: c.KubernetesTopologyZone,
KubernetesEnableEastWest: c.KubernetesEnableEastWest,
KubernetesEastWestDomain: c.KubernetesEastWestDomain,
KubernetesEastWestRangeDomains: c.KubernetesEastWestRangeDomains.values,
Expand Down
1 change: 1 addition & 0 deletions dataclients/kubernetes/clusterclient.go
Original file line number Diff line number Diff line change
Expand Up @@ -677,6 +677,7 @@ func (c *clusterClient) fetchClusterState() (*clusterState, error) {
routeGroups: routeGroups,
services: services,
cachedEndpoints: make(map[endpointID][]string),
cachedEndpointSlices: make(map[endpointID][]skipperEndpoint),
enableEndpointSlices: c.enableEndpointSlices,
}

Expand Down
113 changes: 88 additions & 25 deletions dataclients/kubernetes/clusterstate.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ type clusterState struct {
endpointSlices map[definitions.ResourceID]*skipperEndpointSlice
secrets map[definitions.ResourceID]*secret
cachedEndpoints map[endpointID][]string
cachedEndpointSlices map[endpointID][]skipperEndpoint
enableEndpointSlices bool
}

Expand Down Expand Up @@ -49,8 +50,8 @@ func (state *clusterState) getServiceRG(namespace, name string) (*service, error
return s, nil
}

// GetEndpointsByService returns the skipper endpoints for kubernetes endpoints or endpointslices.
func (state *clusterState) GetEndpointsByService(zone, namespace, name, protocol string, servicePort *servicePort) []string {
// GetEndpointsByService returns the skipper endpoints for kubernetes endpoints.
Comment thread
a4180p marked this conversation as resolved.
func (state *clusterState) GetEndpointsByService(namespace, name, protocol string, servicePort *servicePort) []string {
Comment thread
greeshma1196 marked this conversation as resolved.
epID := endpointID{
ResourceID: newResourceID(namespace, name),
Protocol: protocol,
Expand All @@ -64,23 +65,47 @@ func (state *clusterState) GetEndpointsByService(zone, namespace, name, protocol
}

var targets []string
if state.enableEndpointSlices {
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
targets = eps.targetsByServicePort(zone, "TCP", protocol, servicePort)
} else {
return nil
}
var ep *endpoint
ep, ok := state.endpoints[epID.ResourceID]
if !ok {
return nil
}

targets = ep.targetsByServicePort(protocol, servicePort)

sort.Strings(targets)
state.cachedEndpoints[epID] = targets
Comment thread
MustafaSaber marked this conversation as resolved.
return targets
}

// GetEndpointSlicesByService returns the skipper endpointslices for kubernetes endpointslices.
func (state *clusterState) GetEndpointSlicesByService(zone, namespace, name, protocol string, servicePort *servicePort) []skipperEndpoint {
epID := endpointID{
ResourceID: newResourceID(namespace, name),
Protocol: protocol,
TargetPort: servicePort.TargetPort.String(),
}

state.mu.Lock()
defer state.mu.Unlock()
var targets []skipperEndpoint
if cached, ok := state.cachedEndpointSlices[epID]; ok {
targets = cached
} else {
if ep, ok := state.endpoints[epID.ResourceID]; ok {
targets = ep.targetsByServicePort(protocol, servicePort)
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
targets = eps.targetsByServicePort("TCP", protocol, servicePort)
} else {
return nil
}

sort.Slice(targets, func(i, j int) bool {
return targets[i].Address < targets[j].Address
})
Comment thread
a4180p marked this conversation as resolved.

state.cachedEndpointSlices[epID] = targets
}

sort.Strings(targets)
state.cachedEndpoints[epID] = targets
return targets
return filterByZone(zone, targets)
}

// getEndpointAddresses returns the list of all addresses for the given service using endpoints or endpointslices.
Expand Down Expand Up @@ -113,8 +138,8 @@ func (state *clusterState) getEndpointAddresses(zone, namespace, name string) []
return addresses
}

// GetEndpointsByTarget returns the skipper endpoints for kubernetes endpoints or endpointslices.
func (state *clusterState) GetEndpointsByTarget(zone, namespace, name, protocol, scheme string, target *definitions.BackendPort) []string {
// GetEndpointsByTarget returns the skipper endpoints for kubernetes endpoints.
func (state *clusterState) GetEndpointsByTarget(namespace, name, protocol, scheme string, target *definitions.BackendPort) []string {
Comment thread
greeshma1196 marked this conversation as resolved.
Comment thread
MustafaSaber marked this conversation as resolved.
epID := endpointID{
ResourceID: newResourceID(namespace, name),
Protocol: protocol,
Expand All @@ -128,21 +153,59 @@ func (state *clusterState) GetEndpointsByTarget(zone, namespace, name, protocol,
}

var targets []string
if state.enableEndpointSlices {
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
targets = eps.targetsByServiceTarget(zone, protocol, scheme, target)
} else {
return nil
}
var ep *endpoint
ep, ok := state.endpoints[epID.ResourceID]
if !ok {
return nil
}

targets = ep.targetsByServiceTarget(scheme, target)

sort.Strings(targets)
state.cachedEndpoints[epID] = targets
return targets
}

// GetEndpointSlicesByTarget returns the skipper endpointslices for kubernetes endpointslices.
func (state *clusterState) GetEndpointSlicesByTarget(zone, namespace, name, protocol, scheme string, target *definitions.BackendPort) []skipperEndpoint {
epID := endpointID{
ResourceID: newResourceID(namespace, name),
Protocol: protocol,
TargetPort: target.String(),
}

state.mu.Lock()
defer state.mu.Unlock()
var targets []skipperEndpoint
if cached, ok := state.cachedEndpointSlices[epID]; ok {
Comment thread
MustafaSaber marked this conversation as resolved.
targets = cached
} else {
if ep, ok := state.endpoints[epID.ResourceID]; ok {
targets = ep.targetsByServiceTarget(scheme, target)
if eps, ok := state.endpointSlices[epID.ResourceID]; ok {
targets = eps.targetsByServiceTarget(protocol, scheme, target)
} else {
return nil
}

sort.Slice(targets, func(i, j int) bool {
return targets[i].Address < targets[j].Address
})
state.cachedEndpointSlices[epID] = targets
}

sort.Strings(targets)
state.cachedEndpoints[epID] = targets
return filterByZone(zone, targets)
}

func filterByZone(zone string, targets []skipperEndpoint) []skipperEndpoint {
if zone != "" {
var zoneTargets []skipperEndpoint
for _, target := range targets {
if target.Zone == zone {
zoneTargets = append(zoneTargets, target)
}
}
if len(zoneTargets) >= minEndpointsByZone {
return zoneTargets
}
}
return targets
}
2 changes: 1 addition & 1 deletion dataclients/kubernetes/clusterstate_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,7 @@ func benchmarkCachedEndpoints(b *testing.B, n int) {
b.ResetTimer()
dummy := []string{}
for i := 0; i < b.N; i++ {
dummy = cs.GetEndpointsByTarget("", "default", "foo-0", "TCP", "http", &definitions.BackendPort{})
dummy = cs.GetEndpointsByTarget("default", "foo-0", "TCP", "http", &definitions.BackendPort{})
}
dummy2 = dummy
}
Expand Down
28 changes: 8 additions & 20 deletions dataclients/kubernetes/endpointslices.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ func (eps *skipperEndpointSlice) getPort(protocol, pName string, pValue int) int

return port
}
func (eps *skipperEndpointSlice) targetsByServicePort(zone, protocol, scheme string, servicePort *servicePort) []string {
func (eps *skipperEndpointSlice) targetsByServicePort(protocol, scheme string, servicePort *servicePort) []skipperEndpoint {
var port int
if servicePort.Name != "" {
port = eps.getPort(protocol, servicePort.Name, servicePort.Port)
Expand All @@ -57,35 +57,23 @@ func (eps *skipperEndpointSlice) targetsByServicePort(zone, protocol, scheme str
port = eps.getPort(protocol, servicePort.Name, servicePort.Port)
}

result := make([]string, 0, len(eps.Endpoints))
resultByZone := make([]string, 0, len(eps.Endpoints))
result := make([]skipperEndpoint, 0, len(eps.Endpoints))
for _, ep := range eps.Endpoints {
if ep.Zone == zone {
resultByZone = append(resultByZone, formatEndpointString(ep.Address, scheme, port))
}
result = append(result, formatEndpointString(ep.Address, scheme, port))
}
if len(resultByZone) >= minEndpointsByZone {
return resultByZone
addr := formatEndpointString(ep.Address, scheme, port)
result = append(result, skipperEndpoint{Address: addr, Zone: ep.Zone})
}
return result
}

func (eps *skipperEndpointSlice) targetsByServiceTarget(zone, protocol, scheme string, serviceTarget *definitions.BackendPort) []string {
func (eps *skipperEndpointSlice) targetsByServiceTarget(protocol, scheme string, serviceTarget *definitions.BackendPort) []skipperEndpoint {
pName, _ := serviceTarget.Value.(string)
pValue, _ := serviceTarget.Value.(int)
port := eps.getPort(protocol, pName, pValue)

result := make([]string, 0, len(eps.Endpoints))
resultByZone := make([]string, 0, len(eps.Endpoints))
result := make([]skipperEndpoint, 0, len(eps.Endpoints))
for _, ep := range eps.Endpoints {
if ep.Zone == zone {
resultByZone = append(resultByZone, formatEndpointString(ep.Address, scheme, port))
}
result = append(result, formatEndpointString(ep.Address, scheme, port))
}
if len(resultByZone) >= minEndpointsByZone {
return resultByZone
addr := formatEndpointString(ep.Address, scheme, port)
result = append(result, skipperEndpoint{Address: addr, Zone: ep.Zone})
}
return result
}
Expand Down
85 changes: 63 additions & 22 deletions dataclients/kubernetes/ingressv1.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,11 +82,14 @@ func convertPathRuleV1(
}

var (
eps []string
err error
svc *service
eps []string
epSlices []skipperEndpoint
err error
svc *service
)

dataclientZone := ic.zone

var hostRegexp []string
if host != "" {
hostRegexp = []string{createHostRx(host)}
Expand Down Expand Up @@ -121,7 +124,14 @@ func convertPathRuleV1(
protocol = p
}

eps = state.GetEndpointsByService(ic.zone, ns, svcName, protocol, servicePort)
if state.enableEndpointSlices {

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Kubernetes deprecated endpoints in favor of endpointslices and we used this switch only to enable it when it was more recent.
I think we can get rid of this.
Let's do the cleanup in a different PR.

epSlices = state.GetEndpointSlicesByService(dataclientZone, ns, svcName, protocol, servicePort)
for _, ep := range epSlices {
Comment thread
MustafaSaber marked this conversation as resolved.
eps = append(eps, ep.Address)
}
} else {
eps = state.GetEndpointsByService(ns, svcName, protocol, servicePort)
}
}
if len(eps) == 0 {
ic.logger.Tracef("Target endpoints not found, shuntroute for %s:%s", svcName, svcPort)
Expand All @@ -138,6 +148,7 @@ func convertPathRuleV1(
}

ic.logger.Tracef("Found %d endpoints for %s:%s", len(eps), svcName, svcPort)

if len(eps) == 1 {
r := &eskip.Route{
Id: routeID(ns, name, host, prule.Path, svcName),
Expand All @@ -154,10 +165,20 @@ func convertPathRuleV1(
r := &eskip.Route{
Id: routeID(ns, name, host, prule.Path, svcName),
BackendType: eskip.LBBackend,
LBEndpoints: eps,
LBAlgorithm: getLoadBalancerAlgorithm(metadata, defaultLoadBalancerAlgorithm),
HostRegexps: hostRegexp,
}

if state.enableEndpointSlices {
var lbeps []*eskip.LBEndpoint
for _, ep := range epSlices {
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep.Address, Zone: ep.Zone})
}
r.LBEndpoints = lbeps
} else {
r.LBEndpoints = eskip.NewLBEndpoints(eps)
}

setPathV1(pathMode, r, prule.PathType, prule.Path)
traffic.apply(r)
return r, nil
Expand Down Expand Up @@ -348,14 +369,17 @@ func (ing *ingress) convertDefaultBackendV1(
}

var (
eps []string
err error
ns = i.Metadata.Namespace
name = i.Metadata.Name
svcName = i.Spec.DefaultBackend.Service.Name
svcPort = i.Spec.DefaultBackend.Service.Port
eps []string
epSlices []skipperEndpoint
err error
ns = i.Metadata.Namespace
name = i.Metadata.Name
svcName = i.Spec.DefaultBackend.Service.Name
svcPort = i.Spec.DefaultBackend.Service.Port
)

dataclientZone := ic.zone

svc, err := state.getService(ns, svcName)
if err != nil {
ic.logger.Errorf("Failed to get service %s, %s", svcName, svcPort)
Expand All @@ -382,14 +406,20 @@ func (ing *ingress) convertDefaultBackendV1(
protocol = p
}

eps = state.GetEndpointsByService(
ing.zone,
ns,
svcName,
protocol,
servicePort,
)
ic.logger.Debugf("Found %d endpoints for %s: %v", len(eps), svcName, err)
if state.enableEndpointSlices {
epSlices = state.GetEndpointSlicesByService(dataclientZone, ns, svcName, protocol, servicePort)
for _, ep := range epSlices {
eps = append(eps, ep.Address)
}
} else {
eps = state.GetEndpointsByService(
ns,
svcName,
protocol,
servicePort,
)
ic.logger.Debugf("Found %d endpoints for %s: %v", len(eps), svcName, err)
}
}

if len(eps) == 0 {
Expand All @@ -408,12 +438,23 @@ func (ing *ingress) convertDefaultBackendV1(
}, true, nil
}

return &eskip.Route{
r := &eskip.Route{
Id: routeID(ns, name, "", "", ""),
BackendType: eskip.LBBackend,
LBEndpoints: eps,
LBAlgorithm: getLoadBalancerAlgorithm(i.Metadata, ing.defaultLoadBalancerAlgorithm),
}, true, nil
}

if state.enableEndpointSlices {
var lbeps []*eskip.LBEndpoint
for _, ep := range epSlices {
lbeps = append(lbeps, &eskip.LBEndpoint{Address: ep.Address, Zone: ep.Zone})
}
r.LBEndpoints = lbeps
} else {
r.LBEndpoints = eskip.NewLBEndpoints(eps)
}

return r, true, nil
}

func serviceNameBackend(svcName, svcNamespace string, servicePort *servicePort) string {
Expand Down
Loading
Loading