Skip to content
Open
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 @@ -186,6 +186,7 @@ type Config struct {
KubernetesPathModeString string `yaml:"kubernetes-path-mode"`
KubernetesPathMode kubernetes.PathMode `yaml:"-"`
KubernetesNamespace string `yaml:"kubernetes-namespace"`
KubernetesListChunkSize int `yaml:"kubernetes-list-chunk-size"`
KubernetesEnableEndpointSlices bool `yaml:"enable-kubernetes-endpointslices"`
KubernetesTopologyZone string `yaml:"kubernetes-topology-zone"`
KubernetesEnableEastWest bool `yaml:"enable-kubernetes-east-west"`
Expand Down Expand Up @@ -593,6 +594,7 @@ func NewConfig() *Config {
flag.StringVar(&cfg.WhitelistedHealthCheckCIDR, "whitelisted-healthcheck-cidr", "", "sets the iprange/CIDRS to be whitelisted during healthcheck")
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.IntVar(&cfg.KubernetesListChunkSize, "kubernetes-list-chunk-size", 0, "configures the chunk size (limit parameter) when listing Kubernetes resources. If 0 or negative, pagination is disabled")
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")
Expand Down Expand Up @@ -1110,6 +1112,7 @@ func (c *Config) ToOptions() skipper.Options {
WhitelistedHealthCheckCIDR: whitelistCIDRS,
KubernetesPathMode: c.KubernetesPathMode,
KubernetesNamespace: c.KubernetesNamespace,
KubernetesListChunkSize: c.KubernetesListChunkSize,
KubernetesEnableEndpointslices: c.KubernetesEnableEndpointSlices,
KubernetesTopologyZone: c.KubernetesTopologyZone,
KubernetesEnableEastWest: c.KubernetesEnableEastWest,
Expand Down
100 changes: 76 additions & 24 deletions dataclients/kubernetes/clusterclient.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"os"
"regexp"
"sort"
"strconv"
"strings"
"time"

Expand Down Expand Up @@ -78,6 +79,7 @@ type clusterClient struct {
routeGroupsLabelSelectors string

enableEndpointSlices bool
listChunkSize int

loggedMissingRouteGroups bool
routeGroupValidator *definitions.RouteGroupValidator
Expand Down Expand Up @@ -177,6 +179,7 @@ func newClusterClient(o Options, apiURL, ingCls, rgCls string, quit <-chan struc
routeGroupValidator: &definitions.RouteGroupValidator{EnableAdvancedValidation: false},
ingressValidator: &definitions.IngressV1Validator{EnableAdvancedValidation: false},
enableEndpointSlices: o.KubernetesEnableEndpointslices,
listChunkSize: o.KubernetesListChunkSize,
zone: o.TopologyZone,
ingressStatusFromService: o.IngressStatusFromService,
}
Expand Down Expand Up @@ -327,6 +330,55 @@ func (c *clusterClient) getJSON(uri string, a any) error {
return err
}

type listMetadata struct {
Continue string `json:"continue"`
}

type paginatedList[T any] struct {
Metadata listMetadata `json:"metadata"`
Items []T `json:"items"`
}

func appendQueryParam(uri, key, val string) string {
sep := "?"
if strings.Contains(uri, "?") {
sep = "&"
}
return uri + sep + url.QueryEscape(key) + "=" + url.QueryEscape(val)
}

func loadChunkedList[T any](c *clusterClient, uri string) ([]T, error) {
if c.listChunkSize <= 0 {
var list paginatedList[T]
if err := c.getJSON(uri, &list); err != nil {
return nil, err
}
return list.Items, nil
}

var items []T
continueToken := ""
for {
pageURI := appendQueryParam(uri, "limit", strconv.Itoa(c.listChunkSize))
if continueToken != "" {
pageURI = appendQueryParam(pageURI, "continue", continueToken)
}

var page paginatedList[T]
if err := c.getJSON(pageURI, &page); err != nil {
return nil, err
}

items = append(items, page.Items...)
if page.Metadata.Continue == "" {
break
}
continueToken = page.Metadata.Continue
}

return items, nil
}

func (c *clusterClient) clusterHasRouteGroups() (bool, error) {
var crl ClusterResourceList
if err := c.getJSON(ZalandoResourcesClusterURI, &crl); err != nil { // it probably should bounce once
Expand Down Expand Up @@ -391,14 +443,14 @@ func sortByMetadata(slice any, getMetadata func(int) *definitions.Metadata) {
}

func (c *clusterClient) loadIngressesV1() ([]*definitions.IngressV1Item, error) {
var il definitions.IngressV1List
if err := c.getJSON(c.ingressesURI+c.ingressLabelSelectors, &il); err != nil {
items, err := loadChunkedList[*definitions.IngressV1Item](c, c.ingressesURI+c.ingressLabelSelectors)
if err != nil {
log.Debugf("requesting all ingresses failed: %v", err)
return nil, err
}
log.Debugf("all ingresses received: %d", len(il.Items))
log.Debugf("all ingresses received: %d", len(items))

fItems := c.filterIngressesV1ByClass(il.Items)
fItems := c.filterIngressesV1ByClass(items)
log.Debugf("filtered ingresses by ingress class: %d", len(fItems))

sortByMetadata(fItems, func(i int) *definitions.Metadata { return fItems[i].Metadata })
Expand All @@ -416,12 +468,12 @@ func (c *clusterClient) loadIngressesV1() ([]*definitions.IngressV1Item, error)
}

func (c *clusterClient) LoadRouteGroups() ([]*definitions.RouteGroupItem, error) {
var rgl definitions.RouteGroupList
if err := c.getJSON(c.routeGroupsURI+c.routeGroupsLabelSelectors, &rgl); err != nil {
items, err := loadChunkedList[*definitions.RouteGroupItem](c, c.routeGroupsURI+c.routeGroupsLabelSelectors)
if err != nil {
return nil, err
}
log.Debugf("all routegroups received: %d", len(rgl.Items))
rgl = definitions.NewRouteGroupListWithSharedCache(rgl.Items)
log.Debugf("all routegroups received: %d", len(items))
rgl := definitions.NewRouteGroupListWithSharedCache(items)

rgs := make([]*definitions.RouteGroupItem, 0, len(rgl.Items))
for _, i := range rgl.Items {
Expand Down Expand Up @@ -451,16 +503,16 @@ func (c *clusterClient) LoadRouteGroups() ([]*definitions.RouteGroupItem, error)
}

func (c *clusterClient) loadServices() (map[definitions.ResourceID]*service, error) {
var services serviceList
if err := c.getJSON(c.servicesURI+c.servicesLabelSelectors, &services); err != nil {
items, err := loadChunkedList[*service](c, c.servicesURI+c.servicesLabelSelectors)
if err != nil {
log.Debugf("requesting all services failed: %v", err)
return nil, err
}
log.Debugf("all services received: %d", len(services.Items))
log.Debugf("all services received: %d", len(items))

result := make(map[definitions.ResourceID]*service)
var hasInvalidService bool
for _, service := range services.Items {
for _, service := range items {
if service == nil || service.Meta == nil || service.Spec == nil {
hasInvalidService = true
continue
Expand All @@ -477,15 +529,15 @@ func (c *clusterClient) loadServices() (map[definitions.ResourceID]*service, err
}

func (c *clusterClient) loadSecrets() (map[definitions.ResourceID]*secret, error) {
var secrets secretList
if err := c.getJSON(c.secretsURI+c.secretsLabelSelectors, &secrets); err != nil {
items, err := loadChunkedList[*secret](c, c.secretsURI+c.secretsLabelSelectors)
if err != nil {
log.Debugf("requesting all secrets failed: %v", err)
return nil, err
}
log.Debugf("all secrets received: %d", len(secrets.Items))
log.Debugf("all secrets received: %d", len(items))

result := make(map[definitions.ResourceID]*secret)
for _, secret := range secrets.Items {
for _, secret := range items {
if secret == nil || secret.Metadata == nil {
continue
}
Expand All @@ -497,15 +549,15 @@ func (c *clusterClient) loadSecrets() (map[definitions.ResourceID]*secret, error
}

func (c *clusterClient) loadEndpoints() (map[definitions.ResourceID]*endpoint, error) {
var endpoints endpointList
if err := c.getJSON(c.endpointsURI+c.endpointsLabelSelectors, &endpoints); err != nil {
items, err := loadChunkedList[*endpoint](c, c.endpointsURI+c.endpointsLabelSelectors)
if err != nil {
log.Debugf("requesting all endpoints failed: %v", err)
return nil, err
}
log.Debugf("all endpoints received: %d", len(endpoints.Items))
log.Debugf("all endpoints received: %d", len(items))

result := make(map[definitions.ResourceID]*endpoint)
for _, endpoint := range endpoints.Items {
for _, endpoint := range items {
resID := endpoint.Meta.ToResourceID()
result[resID] = endpoint
}
Expand All @@ -523,14 +575,14 @@ func (c *clusterClient) loadEndpoints() (map[definitions.ResourceID]*endpoint, e
// non-terminating endpoints that should be in the load balancer of a
// given service, check [endpointSlice.ToResourceID].
func (c *clusterClient) loadEndpointSlices() (map[definitions.ResourceID]*skipperEndpointSlice, error) {
var endpointSlices endpointSliceList
if err := c.getJSON(c.endpointSlicesURI+c.endpointSlicesLabelSelectors, &endpointSlices); err != nil {
items, err := loadChunkedList[*endpointSlice](c, c.endpointSlicesURI+c.endpointSlicesLabelSelectors)
if err != nil {
log.Debugf("requesting all endpointslices failed: %v", err)
return nil, err
}
log.Debugf("all endpointslices received: %d", len(endpointSlices.Items))
log.Debugf("all endpointslices received: %d", len(items))

return collectReadyEndpoints(&endpointSlices), nil
return collectReadyEndpoints(&endpointSliceList{Items: items}), nil
}

func collectReadyEndpoints(endpointSlices *endpointSliceList) map[definitions.ResourceID]*skipperEndpointSlice {
Expand Down
Loading
Loading