diff --git a/pkg/common/helpers/event_recorder.go b/pkg/common/helpers/event_recorder.go index c3f867a53..3c2a89208 100644 --- a/pkg/common/helpers/event_recorder.go +++ b/pkg/common/helpers/event_recorder.go @@ -4,14 +4,14 @@ import ( "context" "k8s.io/apimachinery/pkg/runtime" - "k8s.io/client-go/kubernetes" + eventsv1 "k8s.io/client-go/kubernetes/typed/events/v1" kevents "k8s.io/client-go/tools/events" ) // NewEventRecorder creates a new event recorder for the given controller, it will also log the events func NewEventRecorder(ctx context.Context, scheme *runtime.Scheme, - kubeClient kubernetes.Interface, controllerName string) (kevents.EventRecorder, error) { - broadcaster := kevents.NewBroadcaster(&kevents.EventSinkImpl{Interface: kubeClient.EventsV1()}) + eventsClient eventsv1.EventsV1Interface, controllerName string) (kevents.EventRecorder, error) { + broadcaster := kevents.NewBroadcaster(&kevents.EventSinkImpl{Interface: eventsClient}) err := broadcaster.StartRecordingToSinkWithContext(ctx) if err != nil { return nil, nil diff --git a/pkg/common/helpers/event_recorder_test.go b/pkg/common/helpers/event_recorder_test.go index 7e92eb95b..ae138e6a2 100644 --- a/pkg/common/helpers/event_recorder_test.go +++ b/pkg/common/helpers/event_recorder_test.go @@ -52,8 +52,8 @@ func TestNewEventRecorder(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { ctx := context.TODO() - kubeClient := fakekube.NewSimpleClientset() - recorder, err := NewEventRecorder(ctx, tt.scheme, kubeClient, "test") + kubeClient := fakekube.NewClientset() + recorder, err := NewEventRecorder(ctx, tt.scheme, kubeClient.EventsV1(), "test") if err != nil { t.Errorf("NewEventRecorder() error = %v", err) return diff --git a/pkg/placement/controllers/manager.go b/pkg/placement/controllers/manager.go index b143db0af..72baab258 100644 --- a/pkg/placement/controllers/manager.go +++ b/pkg/placement/controllers/manager.go @@ -44,7 +44,7 @@ func RunControllerManagerWithInformers( clusterClient clusterclient.Interface, clusterInformers clusterinformers.SharedInformerFactory, ) error { - recorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, kubeClient, "placement-controller") + recorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, kubeClient.EventsV1(), "placement-controller") if err != nil { return err } diff --git a/pkg/registration/hub/lease/controller_test.go b/pkg/registration/hub/lease/controller_test.go index 8d1479013..443d84853 100644 --- a/pkg/registration/hub/lease/controller_test.go +++ b/pkg/registration/hub/lease/controller_test.go @@ -171,7 +171,7 @@ func TestSync(t *testing.T) { ctx := context.TODO() syncCtx := testingcommon.NewFakeSyncContext(t, testinghelpers.TestManagedClusterName) - mcEventRecorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, hubClient, "test") + mcEventRecorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, hubClient.EventsV1(), "test") if err != nil { t.Fatal(err) } diff --git a/pkg/registration/hub/manager.go b/pkg/registration/hub/manager.go index 89a16d83e..55bd9502d 100644 --- a/pkg/registration/hub/manager.go +++ b/pkg/registration/hub/manager.go @@ -209,7 +209,7 @@ func (m *HubManagerOptions) RunControllerManagerWithInformers( controllerContext.EventRecorder, ) - mcRecorder, err := commonhelpers.NewEventRecorder(ctx, clusterscheme.Scheme, kubeClient, "registration-controller") + mcRecorder, err := commonhelpers.NewEventRecorder(ctx, clusterscheme.Scheme, kubeClient.EventsV1(), "registration-controller") if err != nil { return err } diff --git a/pkg/registration/register/aws_irsa/aws.go b/pkg/registration/register/aws_irsa/aws.go index 867a397a4..a8fe93e54 100644 --- a/pkg/registration/register/aws_irsa/aws.go +++ b/pkg/registration/register/aws_irsa/aws.go @@ -11,7 +11,7 @@ import ( cluster "open-cluster-management.io/api/client/cluster/clientset/versioned" managedclusterv1client "open-cluster-management.io/api/client/cluster/clientset/versioned/typed/cluster/v1" - managedclusterinformers "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster" + clusterv1informer "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster/v1" managedclusterv1lister "open-cluster-management.io/api/client/cluster/listers/cluster/v1" v1 "open-cluster-management.io/api/cluster/v1" ) @@ -72,10 +72,12 @@ func (v *v1AWSIRSAControl) get(name string) (metav1.Object, error) { return managedcluster, nil } -func NewAWSIRSAControl(hubManagedClusterInformer managedclusterinformers.Interface, hubManagedClusterClient cluster.Interface) (AWSIRSAControl, error) { +func NewAWSIRSAControl( + hubManagedClusterInformer clusterv1informer.ManagedClusterInformer, + hubManagedClusterClient cluster.Interface) (AWSIRSAControl, error) { return &v1AWSIRSAControl{ - hubManagedClusterInformer: hubManagedClusterInformer.V1().ManagedClusters().Informer(), - hubManagedClusterLister: hubManagedClusterInformer.V1().ManagedClusters().Lister(), + hubManagedClusterInformer: hubManagedClusterInformer.Informer(), + hubManagedClusterLister: hubManagedClusterInformer.Lister(), hubManagedClusterClient: hubManagedClusterClient.ClusterV1().ManagedClusters(), }, nil } diff --git a/pkg/registration/register/aws_irsa/aws_irsa.go b/pkg/registration/register/aws_irsa/aws_irsa.go index 8916905cc..a3c3f2748 100644 --- a/pkg/registration/register/aws_irsa/aws_irsa.go +++ b/pkg/registration/register/aws_irsa/aws_irsa.go @@ -101,7 +101,7 @@ func (c *AWSIRSADriver) BuildClients(_ context.Context, secretOption register.Se if err != nil { return nil, err } - c.awsIRSAControl, err = NewAWSIRSAControl(clients.ClusterInfomerFactory.Cluster(), clients.ClusterClient) + c.awsIRSAControl, err = NewAWSIRSAControl(clients.ClusterInformer, clients.ClusterClient) if err != nil { return nil, fmt.Errorf("failed to create AWS IRSA control: %w", err) } diff --git a/pkg/registration/register/common.go b/pkg/registration/register/common.go index 17e98cb22..ed1c35396 100644 --- a/pkg/registration/register/common.go +++ b/pkg/registration/register/common.go @@ -13,8 +13,9 @@ import ( "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/util/errors" "k8s.io/apimachinery/pkg/util/wait" - "k8s.io/client-go/informers" "k8s.io/client-go/kubernetes" + leasev1client "k8s.io/client-go/kubernetes/typed/coordination/v1" + eventsv1 "k8s.io/client-go/kubernetes/typed/events/v1" "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" clientcmdapi "k8s.io/client-go/tools/clientcmd/api" @@ -22,9 +23,11 @@ import ( addonclient "open-cluster-management.io/api/client/addon/clientset/versioned" addoninformers "open-cluster-management.io/api/client/addon/informers/externalversions" + addonv1alpha1informers "open-cluster-management.io/api/client/addon/informers/externalversions/addon/v1alpha1" clusterv1client "open-cluster-management.io/api/client/cluster/clientset/versioned" hubclusterclientset "open-cluster-management.io/api/client/cluster/clientset/versioned" - clusterv1informers "open-cluster-management.io/api/client/cluster/informers/externalversions" + clusterinformers "open-cluster-management.io/api/client/cluster/informers/externalversions" + clusterv1informer "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster/v1" clusterv1listers "open-cluster-management.io/api/client/cluster/listers/cluster/v1" clusterv1 "open-cluster-management.io/api/cluster/v1" "open-cluster-management.io/sdk-go/pkg/patcher" @@ -243,15 +246,15 @@ func (a *AggregatedHubDriver) Cleanup(ctx context.Context, cluster *clusterv1.Ma // Clients hold all client needed to connect to hub type Clients struct { - ClusterClient clusterv1client.Interface - KubeClient kubernetes.Interface - AddonClient addonclient.Interface - ClusterInfomerFactory clusterv1informers.SharedInformerFactory - KubeInformerFactory informers.SharedInformerFactory - AddonInformerFactory addoninformers.SharedInformerFactory + ClusterClient clusterv1client.Interface + LeaseClient leasev1client.LeaseInterface + AddonClient addonclient.Interface + EventsClient eventsv1.EventsV1Interface + ClusterInformer clusterv1informer.ManagedClusterInformer + AddonInformer addonv1alpha1informers.ManagedClusterAddOnInformer } -func BuildClientsFromSecretOption(s SecretOption, bootstrap bool) (*Clients, error) { +func KubeConfigFromSecretOption(s SecretOption, bootstrap bool) (*rest.Config, error) { var kubeConfig *rest.Config var err error if bootstrap { @@ -269,12 +272,17 @@ func BuildClientsFromSecretOption(s SecretOption, bootstrap bool) (*Clients, err return nil, fmt.Errorf("unable to load hub kubeconfig from file %q: %w", s.HubKubeconfigFile, err) } } + return kubeConfig, nil +} +func BuildClientsFromConfig(kubeConfig *rest.Config, clusterName string) (*Clients, error) { clients := &Clients{} - clients.KubeClient, err = kubernetes.NewForConfig(kubeConfig) + kubeClient, err := kubernetes.NewForConfig(kubeConfig) if err != nil { return nil, err } + clients.LeaseClient = kubeClient.CoordinationV1().Leases(clusterName) + clients.EventsClient = kubeClient.EventsV1() clients.ClusterClient, err = clusterv1client.NewForConfig(kubeConfig) if err != nil { return nil, err @@ -284,24 +292,25 @@ func BuildClientsFromSecretOption(s SecretOption, bootstrap bool) (*Clients, err return nil, err } - clients.KubeInformerFactory = informers.NewSharedInformerFactoryWithOptions( - clients.KubeClient, - 10*time.Minute, - informers.WithTweakListOptions(func(listOptions *metav1.ListOptions) { - listOptions.LabelSelector = fmt.Sprintf("%s=%s", clusterv1.ClusterNameLabelKey, s.ClusterName) - }), - ) - clients.ClusterInfomerFactory = clusterv1informers.NewSharedInformerFactoryWithOptions( + clients.ClusterInformer = clusterinformers.NewSharedInformerFactoryWithOptions( clients.ClusterClient, 10*time.Minute, - clusterv1informers.WithTweakListOptions(func(listOptions *metav1.ListOptions) { - listOptions.FieldSelector = fields.OneTermEqualSelector("metadata.name", s.ClusterName).String() + clusterinformers.WithTweakListOptions(func(listOptions *metav1.ListOptions) { + listOptions.FieldSelector = fields.OneTermEqualSelector("metadata.name", clusterName).String() }), - ) - clients.AddonInformerFactory = addoninformers.NewSharedInformerFactoryWithOptions( + ).Cluster().V1().ManagedClusters() + clients.AddonInformer = addoninformers.NewSharedInformerFactoryWithOptions( clients.AddonClient, 10*time.Minute, - addoninformers.WithNamespace(s.ClusterName), - ) + addoninformers.WithNamespace(clusterName), + ).Addon().V1alpha1().ManagedClusterAddOns() return clients, nil } + +func BuildClientsFromSecretOption(s SecretOption, bootstrap bool) (*Clients, error) { + kubeConfig, err := KubeConfigFromSecretOption(s, bootstrap) + if err != nil { + return nil, err + } + return BuildClientsFromConfig(kubeConfig, s.ClusterName) +} diff --git a/pkg/registration/register/csr/csr.go b/pkg/registration/register/csr/csr.go index 05a599554..aea991707 100644 --- a/pkg/registration/register/csr/csr.go +++ b/pkg/registration/register/csr/csr.go @@ -19,6 +19,8 @@ import ( "k8s.io/apimachinery/pkg/api/meta" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" utilruntime "k8s.io/apimachinery/pkg/util/runtime" + "k8s.io/client-go/informers" + "k8s.io/client-go/kubernetes" "k8s.io/client-go/tools/cache" clientcmdapi "k8s.io/client-go/tools/clientcmd/api" certutil "k8s.io/client-go/util/cert" @@ -320,11 +322,28 @@ func (c *CSRDriver) Fork(addonName string, secretOption register.SecretOption) r func (c *CSRDriver) BuildClients(ctx context.Context, secretOption register.SecretOption, bootstrap bool) (*register.Clients, error) { logger := klog.FromContext(ctx) - clients, err := register.BuildClientsFromSecretOption(secretOption, bootstrap) + + kubeConfig, err := register.KubeConfigFromSecretOption(secretOption, bootstrap) if err != nil { return nil, err } - csrControl, err := NewCSRControl(logger, clients.KubeInformerFactory.Certificates(), clients.KubeClient) + clients, err := register.BuildClientsFromConfig(kubeConfig, secretOption.ClusterName) + if err != nil { + return nil, err + } + + kubeClient, err := kubernetes.NewForConfig(kubeConfig) + if err != nil { + return nil, err + } + kubeInformerFactory := informers.NewSharedInformerFactoryWithOptions( + kubeClient, + 10*time.Minute, + informers.WithTweakListOptions(func(listOptions *metav1.ListOptions) { + listOptions.LabelSelector = fmt.Sprintf("%s=%s", clusterv1.ClusterNameLabelKey, secretOption.ClusterName) + }), + ) + csrControl, err := NewCSRControl(logger, kubeInformerFactory.Certificates(), kubeClient) if err != nil { return nil, fmt.Errorf("failed to create CSR control: %w", err) } diff --git a/pkg/registration/spoke/lease/lease_controller.go b/pkg/registration/spoke/lease/lease_controller.go index ef7ab9f29..973ced104 100644 --- a/pkg/registration/spoke/lease/lease_controller.go +++ b/pkg/registration/spoke/lease/lease_controller.go @@ -12,7 +12,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" utilruntime "k8s.io/apimachinery/pkg/util/runtime" "k8s.io/apimachinery/pkg/util/wait" - clientset "k8s.io/client-go/kubernetes" + leasev1client "k8s.io/client-go/kubernetes/typed/coordination/v1" clusterv1informer "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster/v1" clusterv1listers "open-cluster-management.io/api/client/cluster/listers/cluster/v1" @@ -32,14 +32,14 @@ type managedClusterLeaseController struct { // NewManagedClusterLeaseController creates a new managed cluster lease controller on the managed cluster. func NewManagedClusterLeaseController( clusterName string, - hubClient clientset.Interface, + leaseClient leasev1client.LeaseInterface, hubClusterInformer clusterv1informer.ManagedClusterInformer, recorder events.Recorder) factory.Controller { c := &managedClusterLeaseController{ clusterName: clusterName, hubClusterLister: hubClusterInformer.Lister(), leaseUpdater: &leaseUpdater{ - hubClient: hubClient, + leaseClient: leaseClient, clusterName: clusterName, leaseName: "managed-cluster-lease", recorder: recorder, @@ -92,7 +92,7 @@ type leaseUpdaterInterface interface { // leaseUpdater periodically updates the lease of a managed cluster type leaseUpdater struct { - hubClient clientset.Interface + leaseClient leasev1client.LeaseInterface clusterName string leaseName string lock sync.Mutex @@ -130,14 +130,14 @@ func (u *leaseUpdater) stop() { // update the lease of a given managed cluster. func (u *leaseUpdater) update(ctx context.Context) { - lease, err := u.hubClient.CoordinationV1().Leases(u.clusterName).Get(ctx, u.leaseName, metav1.GetOptions{}) + lease, err := u.leaseClient.Get(ctx, u.leaseName, metav1.GetOptions{}) if err != nil { utilruntime.HandleError(fmt.Errorf("unable to get cluster lease %q on hub cluster: %w", u.leaseName, err)) return } lease.Spec.RenewTime = &metav1.MicroTime{Time: time.Now()} - if _, err = u.hubClient.CoordinationV1().Leases(u.clusterName).Update(ctx, lease, metav1.UpdateOptions{}); err != nil { + if _, err = u.leaseClient.Update(ctx, lease, metav1.UpdateOptions{}); err != nil { utilruntime.HandleError(fmt.Errorf("unable to update cluster lease %q on hub cluster: %w", u.leaseName, err)) } } diff --git a/pkg/registration/spoke/lease/lease_controller_test.go b/pkg/registration/spoke/lease/lease_controller_test.go index 0c555de64..dd1882417 100644 --- a/pkg/registration/spoke/lease/lease_controller_test.go +++ b/pkg/registration/spoke/lease/lease_controller_test.go @@ -134,7 +134,7 @@ func TestLeaseUpdater(t *testing.T) { initRenewTime := time.Now() hubClient := kubefake.NewSimpleClientset(testinghelpers.NewManagedClusterLease("managed-cluster-lease", initRenewTime)) leaseUpdater := &leaseUpdater{ - hubClient: hubClient, + leaseClient: hubClient.CoordinationV1().Leases(testinghelpers.TestManagedClusterName), clusterName: testinghelpers.TestManagedClusterName, leaseName: "managed-cluster-lease", recorder: eventstesting.NewTestingEventRecorder(t), diff --git a/pkg/registration/spoke/managedcluster/claim_reconcile_test.go b/pkg/registration/spoke/managedcluster/claim_reconcile_test.go index a74bf55d1..07da61c53 100644 --- a/pkg/registration/spoke/managedcluster/claim_reconcile_test.go +++ b/pkg/registration/spoke/managedcluster/claim_reconcile_test.go @@ -113,7 +113,7 @@ func TestSync(t *testing.T) { fakeHubClient := fakekube.NewSimpleClientset() ctx := context.TODO() hubEventRecorder, err := helpers.NewEventRecorder(ctx, - clusterscheme.Scheme, fakeHubClient, "test") + clusterscheme.Scheme, fakeHubClient.EventsV1(), "test") if err != nil { t.Fatal(err) } @@ -344,7 +344,7 @@ func TestExposeClaims(t *testing.T) { fakeHubClient := fakekube.NewSimpleClientset() ctx := context.TODO() hubEventRecorder, err := helpers.NewEventRecorder(ctx, - clusterscheme.Scheme, fakeHubClient, "test") + clusterscheme.Scheme, fakeHubClient.EventsV1(), "test") if err != nil { t.Fatal(err) } diff --git a/pkg/registration/spoke/managedcluster/joining_controller_test.go b/pkg/registration/spoke/managedcluster/joining_controller_test.go index 2e240cf53..bfe25880d 100644 --- a/pkg/registration/spoke/managedcluster/joining_controller_test.go +++ b/pkg/registration/spoke/managedcluster/joining_controller_test.go @@ -84,7 +84,7 @@ func TestSyncManagedCluster(t *testing.T) { fakeHubClient := fakekube.NewSimpleClientset() ctx := context.TODO() hubEventRecorder, err := helpers.NewEventRecorder(ctx, - clusterscheme.Scheme, fakeHubClient, "test") + clusterscheme.Scheme, fakeHubClient.EventsV1(), "test") if err != nil { t.Fatal(err) } diff --git a/pkg/registration/spoke/managedcluster/resource_reconcile_test.go b/pkg/registration/spoke/managedcluster/resource_reconcile_test.go index feb1924ef..1bbb8c15f 100644 --- a/pkg/registration/spoke/managedcluster/resource_reconcile_test.go +++ b/pkg/registration/spoke/managedcluster/resource_reconcile_test.go @@ -318,7 +318,7 @@ func TestHealthCheck(t *testing.T) { ctx := context.TODO() hubEventRecorder, err := helpers.NewEventRecorder(ctx, - clusterscheme.Scheme, fakeHubClient, "test") + clusterscheme.Scheme, fakeHubClient.EventsV1(), "test") if err != nil { t.Fatal(err) } diff --git a/pkg/registration/spoke/registration/hub_timeout_controller.go b/pkg/registration/spoke/registration/hub_timeout_controller.go index aa2eb3afd..2250e5fac 100644 --- a/pkg/registration/spoke/registration/hub_timeout_controller.go +++ b/pkg/registration/spoke/registration/hub_timeout_controller.go @@ -7,7 +7,7 @@ import ( "github.com/openshift/library-go/pkg/controller/factory" "github.com/openshift/library-go/pkg/operator/events" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - clientset "k8s.io/client-go/kubernetes" + leasev1client "k8s.io/client-go/kubernetes/typed/coordination/v1" "k8s.io/klog/v2" ) @@ -17,7 +17,7 @@ const leaseName = "managed-cluster-lease" // TODO: it should finally be part of the lease controller. @xuezhaojun type hubTimeoutController struct { clusterName string - hubClient clientset.Interface + leaseClient leasev1client.LeaseInterface timeoutSeconds int32 lastLeaseRenewTime time.Time handleTimeout func(ctx context.Context) error @@ -25,7 +25,7 @@ type hubTimeoutController struct { func NewHubTimeoutController( clusterName string, - hubClient clientset.Interface, + leaseClient leasev1client.LeaseInterface, timeoutSeconds int32, handleTimeout func(ctx context.Context) error, recorder events.Recorder, @@ -34,7 +34,7 @@ func NewHubTimeoutController( clusterName: clusterName, timeoutSeconds: timeoutSeconds, handleTimeout: handleTimeout, - hubClient: hubClient, + leaseClient: leaseClient, } return factory.New().WithSync(c.sync).ResyncEvery(time.Minute). ToController("HubTimeoutController", recorder) @@ -46,7 +46,7 @@ func (c *hubTimeoutController) sync(ctx context.Context, syncCtx factory.SyncCon return nil } - lease, err := c.hubClient.CoordinationV1().Leases(c.clusterName).Get(ctx, leaseName, metav1.GetOptions{}) + lease, err := c.leaseClient.Get(ctx, leaseName, metav1.GetOptions{}) if err != nil { logger.Error(err, "Failed to get lease", "cluster", c.clusterName, "lease", leaseName) // This means to handle the case which the hub is not connectable from the beginning. diff --git a/pkg/registration/spoke/registration/hub_timeout_controller_test.go b/pkg/registration/spoke/registration/hub_timeout_controller_test.go index 6fb228170..fb6da90ae 100644 --- a/pkg/registration/spoke/registration/hub_timeout_controller_test.go +++ b/pkg/registration/spoke/registration/hub_timeout_controller_test.go @@ -34,13 +34,13 @@ func TestHubTimeoutController_Sync(t *testing.T) { var leaseRenewTime = time.Now() lease := testinghelpers.NewManagedClusterLease("managed-cluster-lease", leaseRenewTime) - leaseClient := kubefake.NewSimpleClientset(lease) + leaseClient := kubefake.NewClientset(lease) time.Sleep(time.Second * time.Duration(c.waitSeconds)) handled := false controller := &hubTimeoutController{ - hubClient: leaseClient, + leaseClient: leaseClient.CoordinationV1().Leases(testinghelpers.TestManagedClusterName), handleTimeout: func(ctx context.Context) error { handled = true return nil diff --git a/pkg/registration/spoke/spokeagent.go b/pkg/registration/spoke/spokeagent.go index 956e3bcd8..d5feb8d08 100644 --- a/pkg/registration/spoke/spokeagent.go +++ b/pkg/registration/spoke/spokeagent.go @@ -224,6 +224,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, if err != nil { return err } + driverInformer, _ := o.driver.InformerHandler() secretInformer := namespacedManagementKubeInformerFactory.Core().V1().Secrets() // Register BootstrapKubeconfigEventHandler as an event handler of secret informer, @@ -282,8 +283,10 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, secretController := register.NewSecretController( secretOption, o.driver, register.GenerateBootstrapStatusUpdater(), recorder, controllerName) - go bootstrapClients.ClusterInfomerFactory.Start(bootstrapCtx.Done()) - go bootstrapClients.KubeInformerFactory.Start(bootstrapCtx.Done()) + go bootstrapClients.ClusterInformer.Informer().Run(bootstrapCtx.Done()) + if driverInformer != nil { + go driverInformer.Run(bootstrapCtx.Done()) + } go spokeClusterCreatingController.Run(bootstrapCtx, 1) go secretController.Run(bootstrapCtx, 1) @@ -315,6 +318,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, if err != nil { return err } + hubDriverInformer, _ := o.driver.InformerHandler() recorder.Event("HubClientConfigReady", "Client config for hub is ready.") // create another RegisterController for registration credential rotation @@ -322,18 +326,18 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, secretController := register.NewSecretController( secretOption, o.driver, register.GenerateStatusUpdater( hubClient.ClusterClient, - hubClient.ClusterInfomerFactory.Cluster().V1().ManagedClusters().Lister(), + hubClient.ClusterInformer.Lister(), o.agentOptions.SpokeClusterName), recorder, controllerName) // create ManagedClusterLeaseController to keep the spoke cluster heartbeat managedClusterLeaseController := lease.NewManagedClusterLeaseController( o.agentOptions.SpokeClusterName, - hubClient.KubeClient, - hubClient.ClusterInfomerFactory.Cluster().V1().ManagedClusters(), + hubClient.LeaseClient, + hubClient.ClusterInformer, recorder, ) - hubEventRecorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, hubClient.KubeClient, "klusterlet-agent") + hubEventRecorder, err := helpers.NewEventRecorder(ctx, clusterscheme.Scheme, hubClient.EventsClient, "klusterlet-agent") if err != nil { return fmt.Errorf("failed to create event recorder: %w", err) } @@ -341,7 +345,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, managedClusterHealthCheckController := managedcluster.NewManagedClusterStatusController( o.agentOptions.SpokeClusterName, hubClient.ClusterClient, - hubClient.ClusterInfomerFactory.Cluster().V1().ManagedClusters(), + hubClient.ClusterInformer, spokeKubeClient.Discovery(), spokeClusterInformerFactory.Cluster().V1alpha1().ClusterClaims(), spokeKubeInformerFactory.Core().V1().Nodes(), @@ -357,7 +361,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, addOnLeaseController = addon.NewManagedClusterAddOnLeaseController( o.agentOptions.SpokeClusterName, hubClient.AddonClient, - hubClient.AddonInformerFactory.Addon().V1alpha1().ManagedClusterAddOns(), + hubClient.AddonInformer, managementKubeClient.CoordinationV1(), spokeKubeClient.CoordinationV1(), AddOnLeaseControllerSyncInterval, //TODO: this interval time should be allowed to change from outside @@ -374,7 +378,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, managementKubeClient, spokeKubeClient, addonDriver, - hubClient.AddonInformerFactory.Addon().V1alpha1().ManagedClusterAddOns(), + hubClient.AddonInformer, recorder, ) } @@ -384,7 +388,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, if features.SpokeMutableFeatureGate.Enabled(ocmfeature.MultipleHubs) { hubAcceptController = registration.NewHubAcceptController( o.agentOptions.SpokeClusterName, - hubClient.ClusterInfomerFactory.Cluster().V1().ManagedClusters(), + hubClient.ClusterInformer, func(ctx context.Context) error { logger.Info("Failed to connect to hub because of hubAcceptClient set to false, restart agent to reselect a new bootstrap kubeconfig") o.agentStopFunc() @@ -395,7 +399,7 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, hubTimeoutController = registration.NewHubTimeoutController( o.agentOptions.SpokeClusterName, - hubClient.KubeClient, + hubClient.LeaseClient, o.registrationOption.HubConnectionTimeoutSeconds, func(ctx context.Context) error { logger.Info("Failed to connect to hub because of lease out-of-date, restart agent to reselect a new bootstrap kubeconfig") @@ -406,10 +410,12 @@ func (o *SpokeAgentConfig) RunSpokeAgentWithSpokeInformers(ctx context.Context, ) } - go hubClient.KubeInformerFactory.Start(ctx.Done()) - go hubClient.ClusterInfomerFactory.Start(ctx.Done()) + if hubDriverInformer != nil { + go hubDriverInformer.Run(ctx.Done()) + } + go hubClient.ClusterInformer.Informer().Run(ctx.Done()) go namespacedManagementKubeInformerFactory.Start(ctx.Done()) - go hubClient.AddonInformerFactory.Start(ctx.Done()) + go hubClient.AddonInformer.Informer().Run(ctx.Done()) go spokeKubeInformerFactory.Start(ctx.Done()) if features.SpokeMutableFeatureGate.Enabled(ocmfeature.ClusterClaim) {