diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml
index 9afab6e0..987c8d09 100644
--- a/.github/workflows/test.yaml
+++ b/.github/workflows/test.yaml
@@ -26,7 +26,7 @@ jobs:
args: --timeout=5m
version: v1.64
- tests:
+ validate:
runs-on: ubuntu-latest
steps:
@@ -40,11 +40,24 @@ jobs:
- name: Validate
run: make validate
+ tests:
+ runs-on: ubuntu-latest
+ needs: validate
+
+ steps:
+ - name: Checkout code
+ uses: actions/checkout@v4
+
+ - uses: actions/setup-go@v5
+ with:
+ go-version-file: go.mod
+
- name: Run unit tests
run: make test-unit
tests-e2e:
runs-on: ubuntu-latest
+ needs: validate
steps:
- name: Checkout code
@@ -56,9 +69,6 @@ jobs:
- uses: actions/setup-go@v5
with:
go-version-file: go.mod
-
- - name: Validate
- run: make validate
- name: Install Ginkgo
run: go install github.com/onsi/ginkgo/v2/ginkgo
diff --git a/charts/k3k/crds/k3k.io_clusters.yaml b/charts/k3k/crds/k3k.io_clusters.yaml
index c6f84fb2..234b14b0 100644
--- a/charts/k3k/crds/k3k.io_clusters.yaml
+++ b/charts/k3k/crds/k3k.io_clusters.yaml
@@ -308,8 +308,6 @@ spec:
In "shared" mode, this also applies to workloads.
type: object
persistence:
- default:
- type: dynamic
description: |-
Persistence specifies options for persisting etcd data.
Defaults to dynamic persistence, which uses a PersistentVolumeClaim to provide data persistence.
@@ -321,6 +319,7 @@ spec:
This field is only relevant in "dynamic" mode.
type: string
storageRequestSize:
+ default: 1G
description: |-
StorageRequestSize is the requested size for the PVC.
This field is only relevant in "dynamic" mode.
@@ -329,8 +328,6 @@ spec:
default: dynamic
description: Type specifies the persistence mode.
type: string
- required:
- - type
type: object
priorityClass:
description: |-
@@ -544,26 +541,6 @@ spec:
description: KubeletPort specefies the port used by k3k-kubelet in
shared mode.
type: integer
- persistence:
- description: Persistence specifies options for persisting etcd data.
- properties:
- storageClassName:
- description: |-
- StorageClassName is the name of the StorageClass to use for the PVC.
- This field is only relevant in "dynamic" mode.
- type: string
- storageRequestSize:
- description: |-
- StorageRequestSize is the requested size for the PVC.
- This field is only relevant in "dynamic" mode.
- type: string
- type:
- default: dynamic
- description: Type specifies the persistence mode.
- type: string
- required:
- - type
- type: object
policyName:
description: PolicyName specifies the virtual cluster policy name
bound to the virtual cluster.
diff --git a/charts/k3k/templates/deployment.yaml b/charts/k3k/templates/deployment.yaml
index cb8f67fa..93eea8b1 100644
--- a/charts/k3k/templates/deployment.yaml
+++ b/charts/k3k/templates/deployment.yaml
@@ -38,6 +38,9 @@ spec:
valueFrom:
fieldRef:
fieldPath: metadata.namespace
+ {{- with .Values.extraEnv }}
+ {{- toYaml . | nindent 10 }}
+ {{- end }}
ports:
- containerPort: 8080
name: https
diff --git a/charts/k3k/values.yaml b/charts/k3k/values.yaml
index fdc2e189..99167297 100644
--- a/charts/k3k/values.yaml
+++ b/charts/k3k/values.yaml
@@ -9,6 +9,19 @@ imagePullSecrets: []
nameOverride: ""
fullnameOverride: ""
+# extraEnv allows you to specify additional environment variables for the k3k controller deployment.
+# This is useful for passing custom configuration or secrets to the controller.
+# For example:
+# extraEnv:
+# - name: MY_CUSTOM_VAR
+# value: "my_custom_value"
+# - name: ANOTHER_VAR
+# valueFrom:
+# secretKeyRef:
+# name: my-secret
+# key: my-key
+extraEnv: []
+
host:
# clusterCIDR specifies the clusterCIDR that will be added to the default networkpolicy, if not set
# the controller will collect the PodCIDRs of all the nodes on the system.
diff --git a/docs/crds/crd-docs.md b/docs/crds/crd-docs.md
index 8ef46852..1da40e59 100644
--- a/docs/crds/crd-docs.md
+++ b/docs/crds/crd-docs.md
@@ -106,7 +106,7 @@ _Appears in:_
| `clusterCIDR` _string_ | ClusterCIDR is the CIDR range for pod IPs.
Defaults to 10.42.0.0/16 in shared mode and 10.52.0.0/16 in virtual mode.
This field is immutable. | | |
| `serviceCIDR` _string_ | ServiceCIDR is the CIDR range for service IPs.
Defaults to 10.43.0.0/16 in shared mode and 10.53.0.0/16 in virtual mode.
This field is immutable. | | |
| `clusterDNS` _string_ | ClusterDNS is the IP address for the CoreDNS service.
Must be within the ServiceCIDR range. Defaults to 10.43.0.10.
This field is immutable. | | |
-| `persistence` _[PersistenceConfig](#persistenceconfig)_ | Persistence specifies options for persisting etcd data.
Defaults to dynamic persistence, which uses a PersistentVolumeClaim to provide data persistence.
A default StorageClass is required for dynamic persistence. | \{ type:dynamic \} | |
+| `persistence` _[PersistenceConfig](#persistenceconfig)_ | Persistence specifies options for persisting etcd data.
Defaults to dynamic persistence, which uses a PersistentVolumeClaim to provide data persistence.
A default StorageClass is required for dynamic persistence. | | |
| `expose` _[ExposeConfig](#exposeconfig)_ | Expose specifies options for exposing the API server.
By default, it's only exposed as a ClusterIP. | | |
| `nodeSelector` _object (keys:string, values:string)_ | NodeSelector specifies node labels to constrain where server/agent pods are scheduled.
In "shared" mode, this also applies to workloads. | | |
| `priorityClass` _string_ | PriorityClass specifies the priorityClassName for server/agent pods.
In "shared" mode, this also applies to workloads. | | |
@@ -203,13 +203,12 @@ PersistenceConfig specifies options for persisting etcd data.
_Appears in:_
- [ClusterSpec](#clusterspec)
-- [ClusterStatus](#clusterstatus)
| Field | Description | Default | Validation |
| --- | --- | --- | --- |
| `type` _[PersistenceMode](#persistencemode)_ | Type specifies the persistence mode. | dynamic | |
| `storageClassName` _string_ | StorageClassName is the name of the StorageClass to use for the PVC.
This field is only relevant in "dynamic" mode. | | |
-| `storageRequestSize` _string_ | StorageRequestSize is the requested size for the PVC.
This field is only relevant in "dynamic" mode. | | |
+| `storageRequestSize` _string_ | StorageRequestSize is the requested size for the PVC.
This field is only relevant in "dynamic" mode. | 1G | |
#### PersistenceMode
diff --git a/main.go b/main.go
index 264cb3d6..aca2993b 100644
--- a/main.go
+++ b/main.go
@@ -35,6 +35,7 @@ var (
k3SImagePullPolicy string
kubeletPortRange string
webhookPortRange string
+ maxConcurrentReconciles int
debug bool
logger *log.Logger
flags = []cli.Flag{
@@ -96,6 +97,13 @@ var (
Usage: "K3K server image pull policy",
Destination: &k3SImagePullPolicy,
},
+ &cli.IntFlag{
+ Name: "max-concurrent-reconciles",
+ EnvVars: []string{"MAX_CONCURRENT_RECONCILES"},
+ Usage: "maximum number of concurrent reconciles",
+ Destination: &maxConcurrentReconciles,
+ Value: 50,
+ },
}
)
@@ -146,19 +154,19 @@ func run(clx *cli.Context) error {
logger.Info("adding cluster controller")
- if err := cluster.Add(ctx, mgr, sharedAgentImage, sharedAgentImagePullPolicy, k3SImage, k3SImagePullPolicy, kubeletPortRange, webhookPortRange); err != nil {
+ if err := cluster.Add(ctx, mgr, sharedAgentImage, sharedAgentImagePullPolicy, k3SImage, k3SImagePullPolicy, kubeletPortRange, webhookPortRange, maxConcurrentReconciles); err != nil {
return fmt.Errorf("failed to add the new cluster controller: %v", err)
}
logger.Info("adding etcd pod controller")
- if err := cluster.AddPodController(ctx, mgr); err != nil {
+ if err := cluster.AddPodController(ctx, mgr, maxConcurrentReconciles); err != nil {
return fmt.Errorf("failed to add the new cluster controller: %v", err)
}
logger.Info("adding clusterpolicy controller")
- if err := policy.Add(mgr, clusterCIDR); err != nil {
+ if err := policy.Add(mgr, clusterCIDR, maxConcurrentReconciles); err != nil {
return fmt.Errorf("failed to add the clusterpolicy controller: %v", err)
}
diff --git a/pkg/apis/k3k.io/v1alpha1/types.go b/pkg/apis/k3k.io/v1alpha1/types.go
index 835eac48..26bbe5fd 100644
--- a/pkg/apis/k3k.io/v1alpha1/types.go
+++ b/pkg/apis/k3k.io/v1alpha1/types.go
@@ -39,7 +39,7 @@ type ClusterSpec struct {
// If not specified, the Kubernetes version of the host node will be used.
//
// +optional
- Version string `json:"version"`
+ Version string `json:"version,omitempty"`
// Mode specifies the cluster provisioning mode: "shared" or "virtual".
// Defaults to "shared". This field is immutable.
@@ -95,8 +95,8 @@ type ClusterSpec struct {
// Defaults to dynamic persistence, which uses a PersistentVolumeClaim to provide data persistence.
// A default StorageClass is required for dynamic persistence.
//
- // +kubebuilder:default={type: "dynamic"}
- Persistence PersistenceConfig `json:"persistence,omitempty"`
+ // +optional
+ Persistence PersistenceConfig `json:"persistence"`
// Expose specifies options for exposing the API server.
// By default, it's only exposed as a ClusterIP.
@@ -120,7 +120,7 @@ type ClusterSpec struct {
// The Secret must have a "token" field in its data.
//
// +optional
- TokenSecretRef *v1.SecretReference `json:"tokenSecretRef"`
+ TokenSecretRef *v1.SecretReference `json:"tokenSecretRef,omitempty"`
// TLSSANs specifies subject alternative names for the K3s server certificate.
//
@@ -212,7 +212,7 @@ type PersistenceConfig struct {
// Type specifies the persistence mode.
//
// +kubebuilder:default="dynamic"
- Type PersistenceMode `json:"type"`
+ Type PersistenceMode `json:"type,omitempty"`
// StorageClassName is the name of the StorageClass to use for the PVC.
// This field is only relevant in "dynamic" mode.
@@ -223,6 +223,7 @@ type PersistenceConfig struct {
// StorageRequestSize is the requested size for the PVC.
// This field is only relevant in "dynamic" mode.
//
+ // +kubebuilder:default="1G"
// +optional
StorageRequestSize string `json:"storageRequestSize,omitempty"`
}
@@ -319,11 +320,6 @@ type ClusterStatus struct {
// +optional
TLSSANs []string `json:"tlsSANs,omitempty"`
- // Persistence specifies options for persisting etcd data.
- //
- // +optional
- Persistence PersistenceConfig `json:"persistence,omitempty"`
-
// PolicyName specifies the virtual cluster policy name bound to the virtual cluster.
//
// +optional
diff --git a/pkg/apis/k3k.io/v1alpha1/zz_generated.deepcopy.go b/pkg/apis/k3k.io/v1alpha1/zz_generated.deepcopy.go
index 5c412243..76fa549c 100644
--- a/pkg/apis/k3k.io/v1alpha1/zz_generated.deepcopy.go
+++ b/pkg/apis/k3k.io/v1alpha1/zz_generated.deepcopy.go
@@ -183,7 +183,6 @@ func (in *ClusterStatus) DeepCopyInto(out *ClusterStatus) {
*out = make([]string, len(*in))
copy(*out, *in)
}
- in.Persistence.DeepCopyInto(&out.Persistence)
}
// DeepCopy is an autogenerated deepcopy function, copying the receiver, creating a new ClusterStatus.
diff --git a/pkg/controller/cluster/cluster.go b/pkg/controller/cluster/cluster.go
index 1cde3eac..fca23c5a 100644
--- a/pkg/controller/cluster/cluster.go
+++ b/pkg/controller/cluster/cluster.go
@@ -32,6 +32,7 @@ import (
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
ctrlruntimeclient "sigs.k8s.io/controller-runtime/pkg/client"
+ ctrlcontroller "sigs.k8s.io/controller-runtime/pkg/controller"
"sigs.k8s.io/controller-runtime/pkg/controller/controllerutil"
"sigs.k8s.io/controller-runtime/pkg/event"
"sigs.k8s.io/controller-runtime/pkg/handler"
@@ -46,14 +47,11 @@ const (
etcdPodFinalizerName = "etcdpod.k3k.io/finalizer"
ClusterInvalidName = "system"
- maxConcurrentReconciles = 1
-
- defaultVirtualClusterCIDR = "10.52.0.0/16"
- defaultVirtualServiceCIDR = "10.53.0.0/16"
- defaultSharedClusterCIDR = "10.42.0.0/16"
- defaultSharedServiceCIDR = "10.43.0.0/16"
- defaultStoragePersistentSize = "1G"
- memberRemovalTimeout = time.Minute * 1
+ defaultVirtualClusterCIDR = "10.52.0.0/16"
+ defaultVirtualServiceCIDR = "10.53.0.0/16"
+ defaultSharedClusterCIDR = "10.42.0.0/16"
+ defaultSharedServiceCIDR = "10.43.0.0/16"
+ memberRemovalTimeout = time.Minute * 1
)
type ClusterReconciler struct {
@@ -68,7 +66,7 @@ type ClusterReconciler struct {
}
// Add adds a new controller to the manager
-func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgentImagePullPolicy, k3SImage string, k3SImagePullPolicy string, kubeletPortRange, webhookPortRange string) error {
+func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgentImagePullPolicy, k3SImage string, k3SImagePullPolicy string, kubeletPortRange, webhookPortRange string, maxConcurrentReconciles int) error {
discoveryClient, err := discovery.NewDiscoveryClientForConfig(mgr.GetConfig())
if err != nil {
return err
@@ -104,6 +102,7 @@ func Add(ctx context.Context, mgr manager.Manager, sharedAgentImage, sharedAgent
Watches(&v1.Namespace{}, namespaceEventHandler(&reconciler)).
Owns(&apps.StatefulSet{}).
Owns(&v1.Service{}).
+ WithOptions(ctrlcontroller.Options{MaxConcurrentReconciles: maxConcurrentReconciles}).
Complete(&reconciler)
}
@@ -227,12 +226,6 @@ func (c *ClusterReconciler) reconcileCluster(ctx context.Context, cluster *v1alp
s := server.New(cluster, c.Client, token, string(cluster.Spec.Mode), c.K3SImage, c.K3SImagePullPolicy)
- cluster.Status.Persistence = cluster.Spec.Persistence
- if cluster.Spec.Persistence.StorageRequestSize == "" {
- // default to 1G of request size
- cluster.Status.Persistence.StorageRequestSize = defaultStoragePersistentSize
- }
-
cluster.Status.ClusterCIDR = cluster.Spec.ClusterCIDR
if cluster.Status.ClusterCIDR == "" {
cluster.Status.ClusterCIDR = defaultVirtualClusterCIDR
diff --git a/pkg/controller/cluster/cluster_suite_test.go b/pkg/controller/cluster/cluster_suite_test.go
index ba93fae0..3cdb2be1 100644
--- a/pkg/controller/cluster/cluster_suite_test.go
+++ b/pkg/controller/cluster/cluster_suite_test.go
@@ -65,7 +65,7 @@ var _ = BeforeSuite(func() {
Expect(err).NotTo(HaveOccurred())
ctx, cancel = context.WithCancel(context.Background())
- err = cluster.Add(ctx, mgr, "rancher/k3k-kubelet:latest", "", "rancher/k3s", "", "50000-51000", "51001-52000")
+ err = cluster.Add(ctx, mgr, "rancher/k3k-kubelet:latest", "", "rancher/k3s", "", "50000-51000", "51001-52000", 50)
Expect(err).NotTo(HaveOccurred())
go func() {
diff --git a/pkg/controller/cluster/cluster_test.go b/pkg/controller/cluster/cluster_test.go
index 1b4ecab9..876e1efe 100644
--- a/pkg/controller/cluster/cluster_test.go
+++ b/pkg/controller/cluster/cluster_test.go
@@ -37,10 +37,8 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
When("creating a Cluster", func() {
- var cluster *v1alpha1.Cluster
-
- BeforeEach(func() {
- cluster = &v1alpha1.Cluster{
+ It("will be created with some defaults", func() {
+ cluster := &v1alpha1.Cluster{
ObjectMeta: metav1.ObjectMeta{
GenerateName: "cluster-",
Namespace: namespace,
@@ -49,15 +47,14 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
err := k8sClient.Create(ctx, cluster)
Expect(err).To(Not(HaveOccurred()))
- })
- It("will be created with some defaults", func() {
Expect(cluster.Spec.Mode).To(Equal(v1alpha1.SharedClusterMode))
Expect(cluster.Spec.Agents).To(Equal(ptr.To[int32](0)))
Expect(cluster.Spec.Servers).To(Equal(ptr.To[int32](1)))
Expect(cluster.Spec.Version).To(BeEmpty())
- // TOFIX
- // Expect(cluster.Spec.Persistence.Type).To(Equal(v1alpha1.DynamicPersistenceMode))
+
+ Expect(cluster.Spec.Persistence.Type).To(Equal(v1alpha1.DynamicPersistenceMode))
+ Expect(cluster.Spec.Persistence.StorageRequestSize).To(Equal("1G"))
serverVersion, err := k8s.DiscoveryClient.ServerVersion()
Expect(err).To(Not(HaveOccurred()))
@@ -94,12 +91,19 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
When("exposing the cluster with nodePort", func() {
It("will have a NodePort service", func() {
- cluster.Spec.Expose = &v1alpha1.ExposeConfig{
- NodePort: &v1alpha1.NodePortConfig{},
+ cluster := &v1alpha1.Cluster{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "cluster-",
+ Namespace: namespace,
+ },
+ Spec: v1alpha1.ClusterSpec{
+ Expose: &v1alpha1.ExposeConfig{
+ NodePort: &v1alpha1.NodePortConfig{},
+ },
+ },
}
- err := k8sClient.Update(ctx, cluster)
- Expect(err).To(Not(HaveOccurred()))
+ Expect(k8sClient.Create(ctx, cluster)).To(Succeed())
var service v1.Service
@@ -119,15 +123,22 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
})
It("will have the specified ports exposed when specified", func() {
- cluster.Spec.Expose = &v1alpha1.ExposeConfig{
- NodePort: &v1alpha1.NodePortConfig{
- ServerPort: ptr.To[int32](30010),
- ETCDPort: ptr.To[int32](30011),
+ cluster := &v1alpha1.Cluster{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "cluster-",
+ Namespace: namespace,
+ },
+ Spec: v1alpha1.ClusterSpec{
+ Expose: &v1alpha1.ExposeConfig{
+ NodePort: &v1alpha1.NodePortConfig{
+ ServerPort: ptr.To[int32](30010),
+ ETCDPort: ptr.To[int32](30011),
+ },
+ },
},
}
- err := k8sClient.Update(ctx, cluster)
- Expect(err).To(Not(HaveOccurred()))
+ Expect(k8sClient.Create(ctx, cluster)).To(Succeed())
var service v1.Service
@@ -161,14 +172,21 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
})
It("will not expose the port when out of range", func() {
- cluster.Spec.Expose = &v1alpha1.ExposeConfig{
- NodePort: &v1alpha1.NodePortConfig{
- ETCDPort: ptr.To[int32](2222),
+ cluster := &v1alpha1.Cluster{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "cluster-",
+ Namespace: namespace,
+ },
+ Spec: v1alpha1.ClusterSpec{
+ Expose: &v1alpha1.ExposeConfig{
+ NodePort: &v1alpha1.NodePortConfig{
+ ETCDPort: ptr.To[int32](2222),
+ },
+ },
},
}
- err := k8sClient.Update(ctx, cluster)
- Expect(err).To(Not(HaveOccurred()))
+ Expect(k8sClient.Create(ctx, cluster)).To(Succeed())
var service v1.Service
@@ -200,12 +218,19 @@ var _ = Describe("Cluster Controller", Label("controller"), Label("Cluster"), fu
When("exposing the cluster with loadbalancer", func() {
It("will have a LoadBalancer service with the default ports exposed", func() {
- cluster.Spec.Expose = &v1alpha1.ExposeConfig{
- LoadBalancer: &v1alpha1.LoadBalancerConfig{},
+ cluster := &v1alpha1.Cluster{
+ ObjectMeta: metav1.ObjectMeta{
+ GenerateName: "cluster-",
+ Namespace: namespace,
+ },
+ Spec: v1alpha1.ClusterSpec{
+ Expose: &v1alpha1.ExposeConfig{
+ LoadBalancer: &v1alpha1.LoadBalancerConfig{},
+ },
+ },
}
- err := k8sClient.Update(ctx, cluster)
- Expect(err).To(Not(HaveOccurred()))
+ Expect(k8sClient.Create(ctx, cluster)).To(Succeed())
var service v1.Service
diff --git a/pkg/controller/cluster/pod.go b/pkg/controller/cluster/pod.go
index 7cd996bf..97b165a8 100644
--- a/pkg/controller/cluster/pod.go
+++ b/pkg/controller/cluster/pod.go
@@ -42,7 +42,7 @@ type PodReconciler struct {
}
// Add adds a new controller to the manager
-func AddPodController(ctx context.Context, mgr manager.Manager) error {
+func AddPodController(ctx context.Context, mgr manager.Manager, maxConcurrentReconciles int) error {
// initialize a new Reconciler
reconciler := PodReconciler{
Client: mgr.GetClient(),
@@ -52,9 +52,7 @@ func AddPodController(ctx context.Context, mgr manager.Manager) error {
return ctrl.NewControllerManagedBy(mgr).
Watches(&v1.Pod{}, handler.EnqueueRequestForOwner(mgr.GetScheme(), mgr.GetRESTMapper(), &apps.StatefulSet{}, handler.OnlyControllerOwner())).
Named(podController).
- WithOptions(controller.Options{
- MaxConcurrentReconciles: maxConcurrentReconciles,
- }).
+ WithOptions(controller.Options{MaxConcurrentReconciles: maxConcurrentReconciles}).
Complete(&reconciler)
}
diff --git a/pkg/controller/cluster/server/server.go b/pkg/controller/cluster/server/server.go
index c123b854..e1f5c4ff 100644
--- a/pkg/controller/cluster/server/server.go
+++ b/pkg/controller/cluster/server/server.go
@@ -379,7 +379,7 @@ func (s *Server) setupDynamicPersistence() v1.PersistentVolumeClaim {
StorageClassName: s.cluster.Spec.Persistence.StorageClassName,
Resources: v1.VolumeResourceRequirements{
Requests: v1.ResourceList{
- "storage": resource.MustParse(s.cluster.Status.Persistence.StorageRequestSize),
+ "storage": resource.MustParse(s.cluster.Spec.Persistence.StorageRequestSize),
},
},
},
diff --git a/pkg/controller/policy/policy.go b/pkg/controller/policy/policy.go
index 40b92a8b..6b9d5344 100644
--- a/pkg/controller/policy/policy.go
+++ b/pkg/controller/policy/policy.go
@@ -16,6 +16,7 @@ import (
"k8s.io/client-go/util/workqueue"
ctrl "sigs.k8s.io/controller-runtime"
"sigs.k8s.io/controller-runtime/pkg/client"
+ "sigs.k8s.io/controller-runtime/pkg/controller"
"sigs.k8s.io/controller-runtime/pkg/event"
"sigs.k8s.io/controller-runtime/pkg/handler"
"sigs.k8s.io/controller-runtime/pkg/manager"
@@ -35,7 +36,7 @@ type VirtualClusterPolicyReconciler struct {
}
// Add the controller to manage the Virtual Cluster policies
-func Add(mgr manager.Manager, clusterCIDR string) error {
+func Add(mgr manager.Manager, clusterCIDR string, maxConcurrentReconciles int) error {
reconciler := VirtualClusterPolicyReconciler{
Client: mgr.GetClient(),
Scheme: mgr.GetScheme(),
@@ -50,6 +51,7 @@ func Add(mgr manager.Manager, clusterCIDR string) error {
Owns(&networkingv1.NetworkPolicy{}).
Owns(&v1.ResourceQuota{}).
Owns(&v1.LimitRange{}).
+ WithOptions(controller.Options{MaxConcurrentReconciles: maxConcurrentReconciles}).
Complete(&reconciler)
}
diff --git a/pkg/controller/policy/policy_suite_test.go b/pkg/controller/policy/policy_suite_test.go
index 42625ccc..03ea70b3 100644
--- a/pkg/controller/policy/policy_suite_test.go
+++ b/pkg/controller/policy/policy_suite_test.go
@@ -54,7 +54,7 @@ var _ = BeforeSuite(func() {
ctrl.SetLogger(zapr.NewLogger(zap.NewNop()))
ctx, cancel = context.WithCancel(context.Background())
- err = policy.Add(mgr, "")
+ err = policy.Add(mgr, "", 50)
Expect(err).NotTo(HaveOccurred())
go func() {