mirror of
https://github.com/open-cluster-management-io/ocm.git
synced 2026-08-29 08:37:21 +00:00
Reduce client scope built from driver (#915)
Scorecard supply-chain security / Scorecard analysis (push) Failing after 1m1s
Post / coverage (push) Failing after 22m36s
Post / images (amd64) (push) Failing after 13m6s
Post / images (arm64) (push) Failing after 2m21s
Post / image manifest (push) Has been skipped
Post / trigger clusteradm e2e (push) Has been skipped
Close stale issues and PRs / stale (push) Successful in 13s
Scorecard supply-chain security / Scorecard analysis (push) Failing after 1m1s
Post / coverage (push) Failing after 22m36s
Post / images (amd64) (push) Failing after 13m6s
Post / images (arm64) (push) Failing after 2m21s
Post / image manifest (push) Has been skipped
Post / trigger clusteradm e2e (push) Has been skipped
Close stale issues and PRs / stale (push) Successful in 13s
Signed-off-by: Jian Qiu <jqiu@redhat.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user