diff --git a/config/config.go b/config/config.go index 2dfd3a91d1..b81483c535 100644 --- a/config/config.go +++ b/config/config.go @@ -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"` @@ -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: ") 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 ..skipper.cluster.local") @@ -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, diff --git a/dataclients/kubernetes/clusterclient.go b/dataclients/kubernetes/clusterclient.go index 4f54befff1..e62642bc85 100644 --- a/dataclients/kubernetes/clusterclient.go +++ b/dataclients/kubernetes/clusterclient.go @@ -14,6 +14,7 @@ import ( "os" "regexp" "sort" + "strconv" "strings" "time" @@ -78,6 +79,7 @@ type clusterClient struct { routeGroupsLabelSelectors string enableEndpointSlices bool + listChunkSize int loggedMissingRouteGroups bool routeGroupValidator *definitions.RouteGroupValidator @@ -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, } @@ -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 @@ -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 }) @@ -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 { @@ -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 @@ -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 } @@ -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 } @@ -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 { diff --git a/dataclients/kubernetes/clusterclient_test.go b/dataclients/kubernetes/clusterclient_test.go index 4883df2b3e..4aeb6b4304 100644 --- a/dataclients/kubernetes/clusterclient_test.go +++ b/dataclients/kubernetes/clusterclient_test.go @@ -2,9 +2,12 @@ package kubernetes_test import ( "bytes" + "encoding/json" + "net/http" "net/http/httptest" "os" "strings" + "sync" "testing" "time" @@ -13,6 +16,7 @@ import ( "github.com/stretchr/testify/require" "github.com/zalando/skipper/dataclients/kubernetes" + "github.com/zalando/skipper/dataclients/kubernetes/definitions" "github.com/zalando/skipper/dataclients/kubernetes/kubernetestest" ) @@ -314,3 +318,183 @@ func TestLoggingInterval(t *testing.T) { assert.Equal(t, 1+n+i, countMessages(), "a new message expected for each subsequent update when log level is debug") } } + +func TestClusterClient_Pagination(t *testing.T) { + type pageResponse struct { + Metadata struct { + Continue string `json:"continue"` + } `json:"metadata"` + Items []definitions.IngressV1Item `json:"items"` + } + + var mu sync.Mutex + var receivedRequests []string + + s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + receivedRequests = append(receivedRequests, r.URL.String()) + mu.Unlock() + + limit := r.URL.Query().Get("limit") + continueToken := r.URL.Query().Get("continue") + + assert.Equal(t, "2", limit) + + var resp pageResponse + switch continueToken { + case "": + resp.Metadata.Continue = "token-1" + resp.Items = []definitions.IngressV1Item{ + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-1"}}, + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-2"}}, + } + case "token-1": + resp.Metadata.Continue = "token-2" + resp.Items = []definitions.IngressV1Item{ + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-3"}}, + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-4"}}, + } + case "token-2": + resp.Metadata.Continue = "" + resp.Items = []definitions.IngressV1Item{ + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-5"}}, + } + default: + http.Error(w, "invalid continue token", http.StatusBadRequest) + return + } + + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(resp) + })) + defer s.Close() + + client, err := kubernetes.NewClusterClient( + kubernetes.Options{ + KubernetesListChunkSize: 2, + }, + s.URL, + "", + "", + nil, + ) + require.NoError(t, err) + + items, err := client.LoadIngressesV1() + require.NoError(t, err) + require.Len(t, items, 5) + + assert.Equal(t, "app-1", items[0].Metadata.Name) + assert.Equal(t, "app-2", items[1].Metadata.Name) + assert.Equal(t, "app-3", items[2].Metadata.Name) + assert.Equal(t, "app-4", items[3].Metadata.Name) + assert.Equal(t, "app-5", items[4].Metadata.Name) + + mu.Lock() + defer mu.Unlock() + require.Len(t, receivedRequests, 3) + assert.Equal(t, "/apis/networking.k8s.io/v1/ingresses?limit=2", receivedRequests[0]) + assert.Equal(t, "/apis/networking.k8s.io/v1/ingresses?limit=2&continue=token-1", receivedRequests[1]) + assert.Equal(t, "/apis/networking.k8s.io/v1/ingresses?limit=2&continue=token-2", receivedRequests[2]) +} + +func TestClusterClient_PaginationDisabled(t *testing.T) { + type pageResponse struct { + Metadata struct { + Continue string `json:"continue"` + } `json:"metadata"` + Items []definitions.IngressV1Item `json:"items"` + } + + var mu sync.Mutex + var receivedRequests []string + + s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + receivedRequests = append(receivedRequests, r.URL.String()) + mu.Unlock() + + assert.Empty(t, r.URL.Query().Get("limit")) + assert.Empty(t, r.URL.Query().Get("continue")) + + var resp pageResponse + resp.Items = []definitions.IngressV1Item{ + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-1"}}, + } + + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(resp) + })) + defer s.Close() + + client, err := kubernetes.NewClusterClient( + kubernetes.Options{ + KubernetesListChunkSize: 0, + }, + s.URL, + "", + "", + nil, + ) + require.NoError(t, err) + + items, err := client.LoadIngressesV1() + require.NoError(t, err) + require.Len(t, items, 1) + + mu.Lock() + defer mu.Unlock() + require.Len(t, receivedRequests, 1) + assert.Equal(t, "/apis/networking.k8s.io/v1/ingresses", receivedRequests[0]) +} + +func TestClusterClient_PaginationWithLabelSelectors(t *testing.T) { + type pageResponse struct { + Metadata struct { + Continue string `json:"continue"` + } `json:"metadata"` + Items []definitions.IngressV1Item `json:"items"` + } + + var mu sync.Mutex + var receivedRequests []string + + s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + mu.Lock() + receivedRequests = append(receivedRequests, r.URL.String()) + mu.Unlock() + + assert.Equal(t, "env=prod", r.URL.Query().Get("labelSelector")) + assert.Equal(t, "1", r.URL.Query().Get("limit")) + + var resp pageResponse + resp.Items = []definitions.IngressV1Item{ + {Metadata: &definitions.Metadata{Namespace: "default", Name: "app-prod"}}, + } + + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(resp) + })) + defer s.Close() + + client, err := kubernetes.NewClusterClient( + kubernetes.Options{ + KubernetesListChunkSize: 1, + IngressLabelSelectors: map[string]string{"env": "prod"}, + }, + s.URL, + "", + "", + nil, + ) + require.NoError(t, err) + + items, err := client.LoadIngressesV1() + require.NoError(t, err) + require.Len(t, items, 1) + + mu.Lock() + defer mu.Unlock() + require.Len(t, receivedRequests, 1) + assert.Equal(t, "/apis/networking.k8s.io/v1/ingresses?labelSelector=env%3Dprod&limit=1", receivedRequests[0]) +} diff --git a/dataclients/kubernetes/export_test.go b/dataclients/kubernetes/export_test.go index 44b94ba955..c9fd9eb8bf 100644 --- a/dataclients/kubernetes/export_test.go +++ b/dataclients/kubernetes/export_test.go @@ -1,7 +1,37 @@ package kubernetes -import "time" +import ( + "time" + + "github.com/zalando/skipper/dataclients/kubernetes/definitions" +) func (c *Client) SetLoggingInterval(d time.Duration) { c.loggingInterval = d } + +type ClusterClient = clusterClient + +func NewClusterClient(o Options, apiURL, ingCls, rgCls string, quit <-chan struct{}) (*clusterClient, error) { + return newClusterClient(o, apiURL, ingCls, rgCls, quit) +} + +func (c *clusterClient) LoadIngressesV1() ([]*definitions.IngressV1Item, error) { + return c.loadIngressesV1() +} + +func (c *clusterClient) LoadServices() (map[definitions.ResourceID]*service, error) { + return c.loadServices() +} + +func (c *clusterClient) LoadEndpoints() (map[definitions.ResourceID]*endpoint, error) { + return c.loadEndpoints() +} + +func (c *clusterClient) LoadSecrets() (map[definitions.ResourceID]*secret, error) { + return c.loadSecrets() +} + +func (c *clusterClient) LoadEndpointSlices() (map[definitions.ResourceID]*skipperEndpointSlice, error) { + return c.loadEndpointSlices() +} diff --git a/dataclients/kubernetes/kube.go b/dataclients/kubernetes/kube.go index 623ec313d0..8a1b64bf2d 100644 --- a/dataclients/kubernetes/kube.go +++ b/dataclients/kubernetes/kube.go @@ -100,6 +100,10 @@ type Options struct { // endpointslices instead of endpoints to scale more than 1000 pods within a service KubernetesEnableEndpointslices bool + // KubernetesListChunkSize configures the chunk size (limit parameter) when listing + // Kubernetes resources. If 0 or negative, pagination is disabled. + KubernetesListChunkSize int + // *DEPRECATED* KubernetesEnableEastWest if set adds automatically routes // with "%s.%s.skipper.cluster.local" domain pattern KubernetesEnableEastWest bool diff --git a/dataclients/kubernetes/kubernetestest/api.go b/dataclients/kubernetes/kubernetestest/api.go index 888e511850..c1781071bc 100644 --- a/dataclients/kubernetes/kubernetestest/api.go +++ b/dataclients/kubernetes/kubernetestest/api.go @@ -7,6 +7,7 @@ import ( "io" "net/http" "regexp" + "strconv" "strings" "sync" @@ -230,7 +231,9 @@ func parseSelectors(r *http.Request) map[string]string { selectors := map[string]string{} for selector := range strings.SplitSeq(rawSelector, ",") { kv := strings.Split(selector, "=") - selectors[kv[0]] = kv[1] + if len(kv) == 2 { + selectors[kv[0]] = kv[1] + } } return selectors @@ -238,7 +241,10 @@ func parseSelectors(r *http.Request) map[string]string { func serve(w http.ResponseWriter, r *http.Request, resources []byte, name string) { selectors := parseSelectors(r) - if name == "" && len(selectors) == 0 { + limitStr := r.URL.Query().Get("limit") + continueToken := r.URL.Query().Get("continue") + + if name == "" && len(selectors) == 0 && limitStr == "" && continueToken == "" { w.Write(resources) return } @@ -296,8 +302,45 @@ func serve(w http.ResponseWriter, r *http.Request, resources []byte, name string } } + var limit int + if limitStr != "" { + var err error + limit, err = strconv.Atoi(limitStr) + if err != nil || limit < 0 { + http.Error(w, "invalid limit parameter", http.StatusBadRequest) + return + } + } + + offset := 0 + if continueToken != "" { + var err error + offset, err = strconv.Atoi(continueToken) + if err != nil || offset < 0 { + http.Error(w, "invalid continue token", http.StatusBadRequest) + return + } + if offset > len(filteredItems) { + w.WriteHeader(http.StatusGone) + return + } + } + + var nextContinue string + if limit > 0 { + end := offset + limit + if end < len(filteredItems) { + nextContinue = strconv.Itoa(end) + filteredItems = filteredItems[offset:end] + } else if offset <= len(filteredItems) { + filteredItems = filteredItems[offset:] + } + } else if offset > 0 { + filteredItems = filteredItems[offset:] + } + var result []byte - if err := itemsJSON(&result, filteredItems); err != nil { + if err := itemsJSONWithContinue(&result, filteredItems, nextContinue); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) } else { w.Write(result) @@ -333,7 +376,16 @@ func initNamespace(kinds map[string][]any) (ns namespace, err error) { } func itemsJSON(b *[]byte, o []any) error { + return itemsJSONWithContinue(b, o, "") +} + +func itemsJSONWithContinue(b *[]byte, o []any, continueToken string) error { items := map[string]any{"items": o} + if continueToken != "" { + items["metadata"] = map[string]any{ + "continue": continueToken, + } + } // converting back to YAML, because we have YAMLToJSON() for bytes, and // the data in `o` contains YAML parser style keys of type interface{} diff --git a/dataclients/kubernetes/kubernetestest/api_test.go b/dataclients/kubernetes/kubernetestest/api_test.go index 51af1ffd51..819adde3a2 100644 --- a/dataclients/kubernetes/kubernetestest/api_test.go +++ b/dataclients/kubernetes/kubernetestest/api_test.go @@ -360,4 +360,43 @@ func TestTestAPI(t *testing.T) { assert.EqualError(t, err, "unexpected status code: 404") }) + + t.Run("pagination chunks", func(t *testing.T) { + // Cluster services has 3 items in total. + var page1 map[string]any + get(t, kubernetes.ServicesClusterURI+"?limit=2", &page1) + check(t, page1, 2, "Service") + + contToken, ok := getField(page1, "metadata", "continue").(string) + require.True(t, ok) + assert.Equal(t, "2", contToken) + + var page2 map[string]any + get(t, kubernetes.ServicesClusterURI+"?limit=2&continue="+contToken, &page2) + check(t, page2, 1, "Service") + + nextCont := getField(page2, "metadata", "continue") + assert.Nil(t, nextCont) + }) + + t.Run("pagination invalid limit", func(t *testing.T) { + var o map[string]any + err := getJSON(s.URL+kubernetes.ServicesClusterURI+"?limit=invalid", &o) + assert.EqualError(t, err, "unexpected status code: 400") + + err = getJSON(s.URL+kubernetes.ServicesClusterURI+"?limit=-1", &o) + assert.EqualError(t, err, "unexpected status code: 400") + }) + + t.Run("pagination invalid continue", func(t *testing.T) { + var o map[string]any + err := getJSON(s.URL+kubernetes.ServicesClusterURI+"?continue=invalid", &o) + assert.EqualError(t, err, "unexpected status code: 400") + }) + + t.Run("pagination expired continue token", func(t *testing.T) { + var o map[string]any + err := getJSON(s.URL+kubernetes.ServicesClusterURI+"?continue=999", &o) + assert.EqualError(t, err, "unexpected status code: 410") + }) } diff --git a/skipper.go b/skipper.go index b5b4c0aa2e..197bfe4553 100644 --- a/skipper.go +++ b/skipper.go @@ -296,6 +296,10 @@ type Options struct { // in the cluster-scope. KubernetesNamespace string + // KubernetesListChunkSize configures the chunk size (limit parameter) when listing + // Kubernetes resources. If 0 or negative, pagination is disabled. + KubernetesListChunkSize int + // KubernetesEnableEndpointslices if set skipper will fetch // endpointslices instead of endpoints to scale more than 1000 // pods within a service @@ -1219,6 +1223,7 @@ func (o *Options) KubernetesDataClientOptions() kubernetes.Options { KubernetesURL: o.KubernetesURL, TokenFile: o.KubernetesTokenFile, KubernetesNamespace: o.KubernetesNamespace, + KubernetesListChunkSize: o.KubernetesListChunkSize, KubernetesEnableEastWest: o.KubernetesEnableEastWest, KubernetesEnableEndpointslices: o.KubernetesEnableEndpointslices, KubernetesEastWestDomain: o.KubernetesEastWestDomain,