Add importer into registration (#753)

* Add importer into registraiton

Signed-off-by: Jian Qiu <jqiu@redhat.com>

* Add unit tests

Signed-off-by: Jian Qiu <jqiu@redhat.com>

* Add integration test

Signed-off-by: Jian Qiu <jqiu@redhat.com>

---------

Signed-off-by: Jian Qiu <jqiu@redhat.com>
This commit is contained in:
Jian Qiu
2024-12-16 13:59:55 +00:00
committed by GitHub
parent 3493630ad2
commit 25ea10bcbf
17 changed files with 3315 additions and 6 deletions
+1 -1
View File
@@ -32,7 +32,7 @@ require (
k8s.io/kube-aggregator v0.31.4
k8s.io/utils v0.0.0-20240921022957-49e7df575cb6
open-cluster-management.io/addon-framework v0.11.1-0.20241129080247-57b1d2859f50
open-cluster-management.io/api v0.15.1-0.20241209025232-b62746ae96d4
open-cluster-management.io/api v0.15.1-0.20241210025410-0ba6809d0ae2
open-cluster-management.io/sdk-go v0.15.1-0.20241125015855-1536c3970f8f
sigs.k8s.io/cluster-inventory-api v0.0.0-20240730014211-ef0154379848
sigs.k8s.io/controller-runtime v0.19.3
+2 -2
View File
@@ -453,8 +453,8 @@ k8s.io/utils v0.0.0-20240921022957-49e7df575cb6 h1:MDF6h2H/h4tbzmtIKTuctcwZmY0tY
k8s.io/utils v0.0.0-20240921022957-49e7df575cb6/go.mod h1:OLgZIPagt7ERELqWJFomSt595RzquPNLL48iOWgYOg0=
open-cluster-management.io/addon-framework v0.11.1-0.20241129080247-57b1d2859f50 h1:TXRd6OdGjArh6cwlCYOqlIcyx21k81oUIYj4rmHlYx0=
open-cluster-management.io/addon-framework v0.11.1-0.20241129080247-57b1d2859f50/go.mod h1:tsBSNs9mGfVQQjXBnjgpiX6r0UM+G3iNfmzQgKhEfw4=
open-cluster-management.io/api v0.15.1-0.20241209025232-b62746ae96d4 h1:f6KU3t9s0PA6vXmAjB6A9sd52OqBqOFK2uAhk3UUBKs=
open-cluster-management.io/api v0.15.1-0.20241209025232-b62746ae96d4/go.mod h1:9erZEWEn4bEqh0nIX2wA7f/s3KCuFycQdBrPrRzi0QM=
open-cluster-management.io/api v0.15.1-0.20241210025410-0ba6809d0ae2 h1:zkp3VJnvexYk5fMf9/yFt6P0fQmp1WFd6Q/Y2t2jF5Q=
open-cluster-management.io/api v0.15.1-0.20241210025410-0ba6809d0ae2/go.mod h1:9erZEWEn4bEqh0nIX2wA7f/s3KCuFycQdBrPrRzi0QM=
open-cluster-management.io/sdk-go v0.15.1-0.20241125015855-1536c3970f8f h1:zeC7QrFNarfK2zY6jGtd+mX+yDrQQmnH/J8A7n5Nh38=
open-cluster-management.io/sdk-go v0.15.1-0.20241125015855-1536c3970f8f/go.mod h1:fi5WBsbC5K3txKb8eRLuP0Sim/Oqz/PHX18skAEyjiA=
sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.30.3 h1:2770sDpzrjjsAtVhSeUFseziht227YAWYHLGNM8QPwY=
+285
View File
@@ -0,0 +1,285 @@
package importer
import (
"context"
"fmt"
"github.com/openshift/api"
"github.com/openshift/library-go/pkg/controller/factory"
"github.com/openshift/library-go/pkg/operator/events"
"github.com/openshift/library-go/pkg/operator/resource/resourceapply"
"github.com/openshift/library-go/pkg/operator/resource/resourcehelper"
"github.com/openshift/library-go/pkg/operator/resource/resourcemerge"
appsv1 "k8s.io/api/apps/v1"
apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
"k8s.io/apimachinery/pkg/api/equality"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/serializer"
utilerrors "k8s.io/apimachinery/pkg/util/errors"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/klog/v2"
"k8s.io/utils/pointer"
clusterclientset "open-cluster-management.io/api/client/cluster/clientset/versioned"
clusterinformerv1 "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster/v1"
clusterlisterv1 "open-cluster-management.io/api/client/cluster/listers/cluster/v1"
operatorclient "open-cluster-management.io/api/client/operator/clientset/versioned"
v1 "open-cluster-management.io/api/cluster/v1"
operatorv1 "open-cluster-management.io/api/operator/v1"
"open-cluster-management.io/sdk-go/pkg/patcher"
"open-cluster-management.io/ocm/pkg/common/queue"
"open-cluster-management.io/ocm/pkg/operator/helpers/chart"
cloudproviders "open-cluster-management.io/ocm/pkg/registration/hub/importer/providers"
)
const (
operatorNamesapce = "open-cluster-management"
bootstrapSA = "cluster-bootstrap"
ManagedClusterConditionImported = "Imported"
)
var (
genericScheme = runtime.NewScheme()
genericCodecs = serializer.NewCodecFactory(genericScheme)
genericCodec = genericCodecs.UniversalDeserializer()
)
func init() {
utilruntime.Must(api.InstallKube(genericScheme))
utilruntime.Must(apiextensionsv1.AddToScheme(genericScheme))
utilruntime.Must(operatorv1.Install(genericScheme))
}
// KlusterletConfigRenderer renders the config for klusterlet chart.
type KlusterletConfigRenderer func(
ctx context.Context, config *chart.KlusterletChartConfig) (*chart.KlusterletChartConfig, error)
type Importer struct {
providers []cloudproviders.Interface
clusterClient clusterclientset.Interface
clusterLister clusterlisterv1.ManagedClusterLister
renders []KlusterletConfigRenderer
patcher patcher.Patcher[*v1.ManagedCluster, v1.ManagedClusterSpec, v1.ManagedClusterStatus]
}
// NewImporter creates an auto import controller
func NewImporter(
renders []KlusterletConfigRenderer,
clusterClient clusterclientset.Interface,
clusterInformer clusterinformerv1.ManagedClusterInformer,
providers []cloudproviders.Interface,
recorder events.Recorder) factory.Controller {
controllerName := "managed-cluster-importer"
syncCtx := factory.NewSyncContext(controllerName, recorder)
i := &Importer{
providers: providers,
clusterClient: clusterClient,
clusterLister: clusterInformer.Lister(),
renders: renders,
patcher: patcher.NewPatcher[
*v1.ManagedCluster, v1.ManagedClusterSpec, v1.ManagedClusterStatus](
clusterClient.ClusterV1().ManagedClusters()),
}
for _, provider := range providers {
provider.Register(syncCtx)
}
return factory.New().WithInformersQueueKeysFunc(queue.QueueKeyByMetaName, clusterInformer.Informer()).
WithSyncContext(syncCtx).WithSync(i.sync).ToController(controllerName, recorder)
}
func (i *Importer) sync(ctx context.Context, syncCtx factory.SyncContext) error {
clusterName := syncCtx.QueueKey()
logger := klog.FromContext(ctx)
logger.V(4).Info("Reconciling key", "clusterName", clusterName)
cluster, err := i.clusterLister.Get(clusterName)
switch {
case errors.IsNotFound(err):
return nil
case err != nil:
return err
}
// If the cluster is imported, skip the reconcile
if meta.IsStatusConditionTrue(cluster.Status.Conditions, ManagedClusterConditionImported) {
return nil
}
// get provider from the provider list
var provider cloudproviders.Interface
for _, p := range i.providers {
if p.IsManagedClusterOwner(cluster) {
provider = p
break
}
}
if provider == nil {
logger.V(2).Info("provider not found for cluster", "cluster", cluster.Name)
return nil
}
newCluster := cluster.DeepCopy()
newCluster, err = i.reconcile(ctx, logger, syncCtx.Recorder(), provider, newCluster)
updated, updatedErr := i.patcher.PatchStatus(ctx, newCluster, newCluster.Status, cluster.Status)
if updatedErr != nil {
return updatedErr
}
if err != nil {
return err
}
if updated {
syncCtx.Recorder().Eventf(
"ManagedClusterImported", "managed cluster %s is imported", clusterName)
}
return nil
}
func (i *Importer) reconcile(
ctx context.Context,
logger klog.Logger,
recorder events.Recorder,
provider cloudproviders.Interface,
cluster *v1.ManagedCluster) (*v1.ManagedCluster, error) {
clients, err := provider.Clients(ctx, cluster)
if err != nil {
meta.SetStatusCondition(&cluster.Status.Conditions, metav1.Condition{
Type: ManagedClusterConditionImported,
Status: metav1.ConditionFalse,
Reason: "KubeConfigGetFailed",
Message: fmt.Sprintf("failed to get kubeconfig. See errors:\n%s",
err.Error()),
})
return cluster, err
}
if clients == nil {
meta.SetStatusCondition(&cluster.Status.Conditions, metav1.Condition{
Type: ManagedClusterConditionImported,
Status: metav1.ConditionFalse,
Reason: "KubeConfigNotFound",
Message: "Secret for kubeconfig is not found.",
})
return cluster, nil
}
// render the klsuterlet chart config
klusterletChartConfig := &chart.KlusterletChartConfig{
CreateNamespace: true,
Klusterlet: chart.KlusterletConfig{
Create: true,
ClusterName: cluster.Name,
ResourceRequirement: operatorv1.ResourceRequirement{
Type: operatorv1.ResourceQosClassDefault,
},
},
}
for _, renderer := range i.renders {
klusterletChartConfig, err = renderer(ctx, klusterletChartConfig)
if err != nil {
meta.SetStatusCondition(&cluster.Status.Conditions, metav1.Condition{
Type: ManagedClusterConditionImported,
Status: metav1.ConditionFalse,
Reason: "ConfigRendererFailed",
Message: fmt.Sprintf("failed to render config. See errors:\n%s",
err.Error()),
})
return cluster, err
}
}
rawManifests, err := chart.RenderKlusterletChart(klusterletChartConfig, operatorNamesapce)
if err != nil {
return cluster, err
}
clientHolder := resourceapply.NewKubeClientHolder(clients.KubeClient).
WithAPIExtensionsClient(clients.APIExtClient).WithDynamicClient(clients.DynamicClient)
cache := resourceapply.NewResourceCache()
var results []resourceapply.ApplyResult
for _, manifest := range rawManifests {
requiredObj, _, err := genericCodec.Decode(manifest, nil, nil)
if err != nil {
logger.Error(err, "failed to decode manifest", "manifest", manifest)
return cluster, err
}
result := resourceapply.ApplyResult{}
switch t := requiredObj.(type) {
case *appsv1.Deployment:
result.Result, result.Changed, result.Error = resourceapply.ApplyDeployment(
ctx, clients.KubeClient.AppsV1(), recorder, t, 0)
results = append(results, result)
case *operatorv1.Klusterlet:
result.Result, result.Changed, result.Error = ApplyKlusterlet(
ctx, clients.OperatorClient, recorder, t)
results = append(results, result)
default:
tempResults := resourceapply.ApplyDirectly(ctx, clientHolder, recorder, cache,
func(name string) ([]byte, error) {
return manifest, nil
},
"manifest")
results = append(results, tempResults...)
}
}
var errs []error
for _, result := range results {
if result.Error != nil {
errs = append(errs, result.Error)
}
}
if len(errs) > 0 {
meta.SetStatusCondition(&cluster.Status.Conditions, metav1.Condition{
Type: ManagedClusterConditionImported,
Status: metav1.ConditionFalse,
Reason: "ImportFailed",
Message: fmt.Sprintf("failed to import the klusterlet. See errors:\n%s",
utilerrors.NewAggregate(errs).Error()),
})
} else {
meta.SetStatusCondition(&cluster.Status.Conditions, metav1.Condition{
Type: ManagedClusterConditionImported,
Status: metav1.ConditionTrue,
Reason: "ImportSucceed",
})
}
return cluster, utilerrors.NewAggregate(errs)
}
func ApplyKlusterlet(
ctx context.Context,
client operatorclient.Interface,
recorder events.Recorder,
required *operatorv1.Klusterlet) (*operatorv1.Klusterlet, bool, error) {
existing, err := client.OperatorV1().Klusterlets().Get(ctx, required.Name, metav1.GetOptions{})
if errors.IsNotFound(err) {
requiredCopy := required.DeepCopy()
actual, err := client.OperatorV1().Klusterlets().Create(ctx, requiredCopy, metav1.CreateOptions{})
resourcehelper.ReportCreateEvent(recorder, required, err)
return actual, true, err
}
if err != nil {
return nil, false, err
}
modified := pointer.Bool(false)
existingCopy := existing.DeepCopy()
resourcemerge.EnsureObjectMeta(modified, &existingCopy.ObjectMeta, required.ObjectMeta)
if !*modified && equality.Semantic.DeepEqual(existingCopy.Spec, required.Spec) {
return existingCopy, false, nil
}
existingCopy.Spec = required.Spec
actual, err := client.OperatorV1().Klusterlets().Update(ctx, existingCopy, metav1.UpdateOptions{})
resourcehelper.ReportUpdateEvent(recorder, required, err)
return actual, true, err
}
@@ -0,0 +1,152 @@
package importer
import (
"context"
"encoding/json"
"testing"
"time"
"github.com/openshift/library-go/pkg/controller/factory"
fakeapiextensions "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset/fake"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
fakedynamic "k8s.io/client-go/dynamic/fake"
kubefake "k8s.io/client-go/kubernetes/fake"
clienttesting "k8s.io/client-go/testing"
fakeclusterclient "open-cluster-management.io/api/client/cluster/clientset/versioned/fake"
clusterinformers "open-cluster-management.io/api/client/cluster/informers/externalversions"
fakeoperatorclient "open-cluster-management.io/api/client/operator/clientset/versioned/fake"
clusterv1 "open-cluster-management.io/api/cluster/v1"
"open-cluster-management.io/sdk-go/pkg/patcher"
testingcommon "open-cluster-management.io/ocm/pkg/common/testing"
"open-cluster-management.io/ocm/pkg/registration/hub/importer/providers"
cloudproviders "open-cluster-management.io/ocm/pkg/registration/hub/importer/providers"
)
func TestSync(t *testing.T) {
cases := []struct {
name string
provider *fakeProvider
key string
cluster *clusterv1.ManagedCluster
validate func(t *testing.T, actions []clienttesting.Action)
}{
{
name: "import succeed",
provider: &fakeProvider{isOwned: true},
key: "cluster1",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
validate: func(t *testing.T, actions []clienttesting.Action) {
testingcommon.AssertActions(t, actions, "patch")
patch := actions[0].(clienttesting.PatchAction).GetPatch()
managedCluster := &clusterv1.ManagedCluster{}
err := json.Unmarshal(patch, managedCluster)
if err != nil {
t.Fatal(err)
}
if !meta.IsStatusConditionTrue(managedCluster.Status.Conditions, ManagedClusterConditionImported) {
t.Errorf("expected managed cluster to be imported")
}
},
},
{
name: "no cluster",
provider: &fakeProvider{isOwned: true},
key: "cluster1",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster2"}},
validate: func(t *testing.T, actions []clienttesting.Action) {
testingcommon.AssertNoActions(t, actions)
},
},
{
name: "not owned by the provider",
provider: &fakeProvider{isOwned: false},
key: "cluster1",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
validate: func(t *testing.T, actions []clienttesting.Action) {
testingcommon.AssertNoActions(t, actions)
},
},
{
name: "clients for remote cluster is not generated",
provider: &fakeProvider{isOwned: true, noClients: true},
key: "cluster1",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
validate: func(t *testing.T, actions []clienttesting.Action) {
testingcommon.AssertActions(t, actions, "patch")
patch := actions[0].(clienttesting.PatchAction).GetPatch()
managedCluster := &clusterv1.ManagedCluster{}
err := json.Unmarshal(patch, managedCluster)
if err != nil {
t.Fatal(err)
}
if !meta.IsStatusConditionFalse(managedCluster.Status.Conditions, ManagedClusterConditionImported) {
t.Errorf("expected managed cluster to be imported")
}
},
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
clusterClient := fakeclusterclient.NewSimpleClientset(c.cluster)
clusterInformer := clusterinformers.NewSharedInformerFactory(
clusterClient, 10*time.Minute).Cluster().V1().ManagedClusters()
clusterStore := clusterInformer.Informer().GetStore()
if err := clusterStore.Add(c.cluster); err != nil {
t.Fatal(err)
}
importer := &Importer{
providers: []cloudproviders.Interface{c.provider},
clusterClient: clusterClient,
clusterLister: clusterInformer.Lister(),
patcher: patcher.NewPatcher[
*clusterv1.ManagedCluster, clusterv1.ManagedClusterSpec, clusterv1.ManagedClusterStatus](
clusterClient.ClusterV1().ManagedClusters()),
}
err := importer.sync(context.TODO(), testingcommon.NewFakeSyncContext(t, c.key))
if err != nil {
t.Fatal(err)
}
c.validate(t, clusterClient.Actions())
})
}
}
type fakeProvider struct {
isOwned bool
noClients bool
kubeConfigErr error
}
// KubeConfig is to return the config to connect to the target cluster.
func (f *fakeProvider) Clients(_ context.Context, _ *clusterv1.ManagedCluster) (*providers.Clients, error) {
if f.kubeConfigErr != nil {
return nil, f.kubeConfigErr
}
if f.noClients {
return nil, nil
}
return &providers.Clients{
KubeClient: kubefake.NewClientset(),
// due to https://github.com/kubernetes/kubernetes/issues/126850, still need to use NewSimpleClientset
APIExtClient: fakeapiextensions.NewSimpleClientset(),
OperatorClient: fakeoperatorclient.NewSimpleClientset(),
DynamicClient: fakedynamic.NewSimpleDynamicClient(runtime.NewScheme()),
}, nil
}
// IsManagedClusterOwner check if the provider is used to manage this cluster
func (f *fakeProvider) IsManagedClusterOwner(_ *clusterv1.ManagedCluster) bool {
return f.isOwned
}
// Register registers the provider to the importer. The provider should enqueue the resource
// into the queue with the name of the managed cluster
func (f *fakeProvider) Register(_ factory.SyncContext) {}
// Run starts the provider
func (f *fakeProvider) Run(_ context.Context) {}
@@ -0,0 +1,17 @@
package options
import "github.com/spf13/pflag"
type Options struct {
APIServerURL string
}
func New() *Options {
return &Options{}
}
// AddFlags registers flags for manager
func (m *Options) AddFlags(fs *pflag.FlagSet) {
fs.StringVar(&m.APIServerURL, "hub-apiserver-url", m.APIServerURL,
"APIServer URL of the hub cluster that the spoke cluster can access, Only used for spoke cluster import")
}
@@ -0,0 +1,161 @@
package capi
import (
"context"
"fmt"
"time"
"github.com/openshift/library-go/pkg/controller/factory"
"github.com/pkg/errors"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/dynamic/dynamicinformer"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
"k8s.io/client-go/tools/clientcmd"
"k8s.io/klog/v2"
clusterinformerv1 "open-cluster-management.io/api/client/cluster/informers/externalversions/cluster/v1"
clusterv1 "open-cluster-management.io/api/cluster/v1"
"open-cluster-management.io/ocm/pkg/registration/hub/importer/providers"
)
var ClusterAPIGVR = schema.GroupVersionResource{
Group: "cluster.x-k8s.io",
Version: "v1beta1",
Resource: "clusters",
}
const (
ByCAPIResource = "by-capi-resource"
CAPIAnnotationKey = "cluster.x-k8s.io/cluster"
)
type CAPIProvider struct {
informer dynamicinformer.DynamicSharedInformerFactory
lister cache.GenericLister
kubeClient kubernetes.Interface
managedClusterIndexer cache.Indexer
}
func NewCAPIProvider(
kubeconfig *rest.Config, clusterInformer clusterinformerv1.ManagedClusterInformer) providers.Interface {
dynamicClient := dynamic.NewForConfigOrDie(kubeconfig)
kubeClient := kubernetes.NewForConfigOrDie(kubeconfig)
dynamicInformer := dynamicinformer.NewDynamicSharedInformerFactory(dynamicClient, 30*time.Minute)
utilruntime.Must(clusterInformer.Informer().AddIndexers(cache.Indexers{
ByCAPIResource: indexByCAPIResource,
}))
return &CAPIProvider{
informer: dynamicInformer,
lister: dynamicInformer.ForResource(ClusterAPIGVR).Lister(),
kubeClient: kubeClient,
managedClusterIndexer: clusterInformer.Informer().GetIndexer(),
}
}
func (c *CAPIProvider) Clients(ctx context.Context, cluster *clusterv1.ManagedCluster) (*providers.Clients, error) {
logger := klog.FromContext(ctx)
clusterKey := capiNameFromManagedCluster(cluster)
namespace, name, err := cache.SplitMetaNamespaceKey(clusterKey)
if err != nil {
return nil, err
}
_, err = c.lister.ByNamespace(namespace).Get(name)
switch {
case apierrors.IsNotFound(err):
logger.V(4).Info("cluster is not found", "name", name, "namespace", namespace)
// TODO(qiujian16) need to consider requeue in this case, since secrets is not watched.
return nil, nil
case err != nil:
return nil, err
}
secret, err := c.kubeClient.CoreV1().Secrets(namespace).Get(ctx, name+"-kubeconfig", metav1.GetOptions{})
switch {
case apierrors.IsNotFound(err):
logger.V(4).Info(
"kubeconfig secret is not found", "name", name+"-kubeconfig", "namespace", namespace)
return nil, nil
case err != nil:
return nil, err
}
data, ok := secret.Data["value"]
if !ok {
return nil, errors.Errorf("missing key %q in secret data", name)
}
configOverride, err := clientcmd.NewClientConfigFromBytes(data)
if err != nil {
return nil, err
}
config, err := configOverride.ClientConfig()
if err != nil {
return nil, err
}
return providers.NewClient(config)
}
func (c *CAPIProvider) Register(syncCtx factory.SyncContext) {
_, err := c.informer.ForResource(ClusterAPIGVR).Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
c.enqueueManagedClusterByCAPI(obj, syncCtx)
},
UpdateFunc: func(oldObj, newObj interface{}) {
c.enqueueManagedClusterByCAPI(newObj, syncCtx)
},
})
utilruntime.HandleError(err)
}
func (c *CAPIProvider) IsManagedClusterOwner(cluster *clusterv1.ManagedCluster) bool {
clusterKey := capiNameFromManagedCluster(cluster)
namespace, name, _ := cache.SplitMetaNamespaceKey(clusterKey)
_, err := c.lister.ByNamespace(namespace).Get(name)
return err == nil
}
func (c *CAPIProvider) Run(ctx context.Context) {
c.informer.Start(ctx.Done())
}
func (c *CAPIProvider) enqueueManagedClusterByCAPI(obj interface{}, syncCtx factory.SyncContext) {
accessor, _ := meta.Accessor(obj)
objs, err := c.managedClusterIndexer.ByIndex(ByCAPIResource, fmt.Sprintf(
"%s/%s", accessor.GetNamespace(), accessor.GetName()))
if err != nil {
return
}
for _, obj := range objs {
accessor, _ := meta.Accessor(obj)
syncCtx.Queue().Add(accessor.GetName())
}
}
func indexByCAPIResource(obj interface{}) ([]string, error) {
cluster, ok := obj.(*clusterv1.ManagedCluster)
if !ok {
return []string{}, nil
}
return []string{capiNameFromManagedCluster(cluster)}, nil
}
func capiNameFromManagedCluster(cluster *clusterv1.ManagedCluster) string {
if len(cluster.Annotations) > 0 {
if key, ok := cluster.Annotations[CAPIAnnotationKey]; ok {
return key
}
}
return fmt.Sprintf("%s/%s", cluster.Name, cluster.Name)
}
@@ -0,0 +1,271 @@
package capi
import (
"context"
"testing"
"github.com/ghodss/yaml"
"github.com/openshift/library-go/pkg/controller/factory"
"github.com/openshift/library-go/pkg/operator/events/eventstesting"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/client-go/dynamic/dynamicinformer"
fakedynamic "k8s.io/client-go/dynamic/fake"
fakekube "k8s.io/client-go/kubernetes/fake"
"k8s.io/client-go/tools/cache"
clientcmdapiv1 "k8s.io/client-go/tools/clientcmd/api/v1"
fakecluster "open-cluster-management.io/api/client/cluster/clientset/versioned/fake"
clusterinformers "open-cluster-management.io/api/client/cluster/informers/externalversions"
clusterv1 "open-cluster-management.io/api/cluster/v1"
testingcommon "open-cluster-management.io/ocm/pkg/common/testing"
)
func TestEnqueu(t *testing.T) {
cases := []struct {
name string
cluster *clusterv1.ManagedCluster
capiName string
capiNamespace string
expectedKey string
}{
{
name: "enqueu by name",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
capiName: "cluster1",
capiNamespace: "cluster1",
expectedKey: "cluster1",
},
{
name: "enqueu by annotation",
cluster: &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster2",
Annotations: map[string]string{
CAPIAnnotationKey: "capi/cluster1",
},
},
},
capiName: "cluster1",
capiNamespace: "capi",
expectedKey: "cluster2",
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
client := fakecluster.NewSimpleClientset(c.cluster)
informerFactory := clusterinformers.NewSharedInformerFactory(client, 0)
clusterInformer := informerFactory.Cluster().V1().ManagedClusters()
if err := clusterInformer.Informer().AddIndexers(cache.Indexers{
ByCAPIResource: indexByCAPIResource,
}); err != nil {
t.Fatal(err)
}
if err := clusterInformer.Informer().GetStore().Add(c.cluster); err != nil {
t.Fatal(err)
}
provider := &CAPIProvider{
managedClusterIndexer: clusterInformer.Informer().GetIndexer(),
}
syncCtx := factory.NewSyncContext("test", eventstesting.NewTestingEventRecorder(t))
provider.enqueueManagedClusterByCAPI(&metav1.PartialObjectMetadata{
ObjectMeta: metav1.ObjectMeta{
Name: c.capiName,
Namespace: c.capiNamespace,
},
}, syncCtx)
if i, _ := syncCtx.Queue().Get(); i.(string) != c.expectedKey {
t.Errorf("expected key %s but got %s", c.expectedKey, syncCtx.QueueKey())
}
})
}
}
func TestClients(t *testing.T) {
cases := []struct {
name string
capiObjects []runtime.Object
kubeObjects []runtime.Object
cluster *clusterv1.ManagedCluster
expectErr bool
}{
{
name: "capi cluster not found",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
},
{
name: "secret not found",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "cluster1", "cluster1")},
},
{
name: "secret found with invalid key",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "cluster1", "cluster1")},
kubeObjects: []runtime.Object{
&corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster1-kubeconfig",
Namespace: "cluster1",
},
},
},
expectErr: true,
},
{
name: "build client successfully",
cluster: &clusterv1.ManagedCluster{ObjectMeta: metav1.ObjectMeta{Name: "cluster1"}},
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "cluster1", "cluster1")},
kubeObjects: []runtime.Object{
func() *corev1.Secret {
clientConfig := clientcmdapiv1.Config{
// Define a cluster stanza based on the bootstrap kubeconfig.
Clusters: []clientcmdapiv1.NamedCluster{
{
Name: "hub",
Cluster: clientcmdapiv1.Cluster{
Server: "https://test",
},
},
},
// Define auth based on the obtained client cert.
AuthInfos: []clientcmdapiv1.NamedAuthInfo{
{
Name: "bootstrap",
AuthInfo: clientcmdapiv1.AuthInfo{
Token: "test",
},
},
},
// Define a context that connects the auth info and cluster, and set it as the default
Contexts: []clientcmdapiv1.NamedContext{
{
Name: "bootstrap",
Context: clientcmdapiv1.Context{
Cluster: "hub",
AuthInfo: "bootstrap",
Namespace: "default",
},
},
},
CurrentContext: "bootstrap",
}
bootstrapConfigBytes, _ := yaml.Marshal(clientConfig)
return &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster1-kubeconfig",
Namespace: "cluster1",
},
Data: map[string][]byte{
"value": bootstrapConfigBytes,
},
}
}(),
},
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
dynamicClient := fakedynamic.NewSimpleDynamicClient(runtime.NewScheme(), c.capiObjects...)
kubeClient := fakekube.NewClientset(c.kubeObjects...)
dynamicInformers := dynamicinformer.NewDynamicSharedInformerFactory(dynamicClient, 0)
for _, capiObj := range c.capiObjects {
if err := dynamicInformers.ForResource(ClusterAPIGVR).Informer().GetStore().Add(capiObj); err != nil {
t.Fatal(err)
}
}
provider := &CAPIProvider{
kubeClient: kubeClient,
informer: dynamicInformers,
lister: dynamicInformers.ForResource(ClusterAPIGVR).Lister(),
}
_, err := provider.Clients(context.TODO(), c.cluster)
if c.expectErr && err == nil {
t.Errorf("expected error but got nil")
}
if !c.expectErr && err != nil {
t.Errorf("expected no error but got %v", err)
}
})
}
}
func TestIsManagedClusterOwner(t *testing.T) {
cases := []struct {
name string
capiObjects []runtime.Object
cluster *clusterv1.ManagedCluster
expectedOwn bool
}{
{
name: "by cluster name",
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "cluster1", "cluster1")},
cluster: &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{Name: "cluster1"},
},
expectedOwn: true,
},
{
name: "by cluster annotation",
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "capi", "cluster1")},
cluster: &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster2",
Annotations: map[string]string{
CAPIAnnotationKey: "capi/cluster1",
},
},
},
expectedOwn: true,
},
{
name: "by cluster annotation",
capiObjects: []runtime.Object{
testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", "capi", "cluster2")},
cluster: &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster2",
Annotations: map[string]string{
CAPIAnnotationKey: "capi/cluster1",
},
},
},
expectedOwn: false,
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
dynamicClient := fakedynamic.NewSimpleDynamicClient(runtime.NewScheme(), c.capiObjects...)
dynamicInformers := dynamicinformer.NewDynamicSharedInformerFactory(dynamicClient, 0)
for _, capiObj := range c.capiObjects {
if err := dynamicInformers.ForResource(ClusterAPIGVR).Informer().GetStore().Add(capiObj); err != nil {
t.Fatal(err)
}
}
provider := &CAPIProvider{
lister: dynamicInformers.ForResource(ClusterAPIGVR).Lister(),
}
owned := provider.IsManagedClusterOwner(c.cluster)
if c.expectedOwn != owned {
t.Errorf("expected owned cluster %t but got %t", c.expectedOwn, owned)
}
})
}
}
@@ -0,0 +1,65 @@
package providers
import (
"context"
"github.com/openshift/library-go/pkg/controller/factory"
apiextensionsclient "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
operatorclient "open-cluster-management.io/api/client/operator/clientset/versioned"
clusterv1 "open-cluster-management.io/api/cluster/v1"
)
// Interface is the interface that a cluster provider should implement
type Interface interface {
// Clients returns the client to connect to the target cluster. The client should have the sufficient
// permission to create CRDs/operator and klusterlet CR in the remote cluster.
Clients(ctx context.Context, cluster *clusterv1.ManagedCluster) (*Clients, error)
// IsManagedClusterOwner check if the provider is used to manage this cluster
IsManagedClusterOwner(cluster *clusterv1.ManagedCluster) bool
// Register registers the provider to the importer. The provider should enqueue the resource
// into the queue with the name of the managed cluster
Register(syncCtx factory.SyncContext)
// Run starts the provider. The provider might need to watch the provider related resources
// on the hub cluster, or start a periodic task.
Run(ctx context.Context)
}
type Clients struct {
KubeClient kubernetes.Interface
APIExtClient apiextensionsclient.Interface
OperatorClient operatorclient.Interface
DynamicClient dynamic.Interface
}
func NewClient(config *rest.Config) (*Clients, error) {
kubeClient, err := kubernetes.NewForConfig(config)
if err != nil {
return nil, err
}
hubApiExtensionClient, err := apiextensionsclient.NewForConfig(config)
if err != nil {
return nil, err
}
operatorClient, err := operatorclient.NewForConfig(config)
if err != nil {
return nil, err
}
dynamicClient, err := dynamic.NewForConfig(config)
if err != nil {
return nil, err
}
return &Clients{
APIExtClient: hubApiExtensionClient,
KubeClient: kubeClient,
OperatorClient: operatorClient,
DynamicClient: dynamicClient,
}, nil
}
@@ -0,0 +1,93 @@
package importer
import (
"context"
"fmt"
"github.com/ghodss/yaml"
authv1 "k8s.io/api/authentication/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/kubernetes"
clientcmdapiv1 "k8s.io/client-go/tools/clientcmd/api/v1"
"k8s.io/utils/ptr"
sdkhelpers "open-cluster-management.io/sdk-go/pkg/helpers"
"open-cluster-management.io/ocm/pkg/operator/helpers/chart"
)
func RenderBootstrapHubKubeConfig(
kubeClient kubernetes.Interface, apiServerURL string) KlusterletConfigRenderer {
return func(ctx context.Context, config *chart.KlusterletChartConfig) (*chart.KlusterletChartConfig, error) {
// get bootstrap token
tr, err := kubeClient.CoreV1().
ServiceAccounts(operatorNamesapce).
CreateToken(ctx, bootstrapSA, &authv1.TokenRequest{
Spec: authv1.TokenRequestSpec{
// token expired in 1 hour
ExpirationSeconds: ptr.To[int64](3600),
},
}, metav1.CreateOptions{})
if err != nil {
return config, fmt.Errorf(
"failed to get token from sa %s/%s: %v", operatorNamesapce, bootstrapSA, err)
}
// get apisever url
url := apiServerURL
if len(url) == 0 {
url, err = sdkhelpers.GetAPIServer(kubeClient)
if err != nil {
return config, err
}
}
// get cabundle
ca, err := sdkhelpers.GetCACert(kubeClient)
if err != nil {
return config, err
}
clientConfig := clientcmdapiv1.Config{
// Define a cluster stanza based on the bootstrap kubeconfig.
Clusters: []clientcmdapiv1.NamedCluster{
{
Name: "hub",
Cluster: clientcmdapiv1.Cluster{
Server: url,
CertificateAuthorityData: ca,
},
},
},
// Define auth based on the obtained client cert.
AuthInfos: []clientcmdapiv1.NamedAuthInfo{
{
Name: "bootstrap",
AuthInfo: clientcmdapiv1.AuthInfo{
Token: tr.Status.Token,
},
},
},
// Define a context that connects the auth info and cluster, and set it as the default
Contexts: []clientcmdapiv1.NamedContext{
{
Name: "bootstrap",
Context: clientcmdapiv1.Context{
Cluster: "hub",
AuthInfo: "bootstrap",
Namespace: "default",
},
},
},
CurrentContext: "bootstrap",
}
bootstrapConfigBytes, err := yaml.Marshal(clientConfig)
if err != nil {
return config, err
}
config.BootstrapHubKubeConfig = string(bootstrapConfigBytes)
return config, nil
}
}
@@ -0,0 +1,99 @@
package importer
import (
"context"
"testing"
"github.com/ghodss/yaml"
authenticationv1 "k8s.io/api/authentication/v1"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
kubefake "k8s.io/client-go/kubernetes/fake"
clienttesting "k8s.io/client-go/testing"
"k8s.io/client-go/tools/clientcmd"
clientcmdapiv1 "k8s.io/client-go/tools/clientcmd/api/v1"
"open-cluster-management.io/ocm/pkg/operator/helpers/chart"
)
func TestRenderBootstrapHubKubeConfig(t *testing.T) {
cases := []struct {
name string
objects []runtime.Object
apiserverURL string
expectedURL string
}{
{
name: "render apiserver from input",
apiserverURL: "https://127.0.0.1:6443",
expectedURL: "https://127.0.0.1:6443",
},
{
name: "render apiserver from cluster-info",
objects: []runtime.Object{
func() *corev1.ConfigMap {
config := clientcmdapiv1.Config{
// Define a cluster stanza based on the bootstrap kubeconfig.
Clusters: []clientcmdapiv1.NamedCluster{
{
Name: "hub",
Cluster: clientcmdapiv1.Cluster{
Server: "https://test",
},
},
},
}
bootstrapConfigBytes, _ := yaml.Marshal(config)
return &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Name: "cluster-info",
Namespace: "kube-public",
},
Data: map[string]string{
"kubeconfig": string(bootstrapConfigBytes),
},
}
}(),
},
expectedURL: "https://test",
},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
client := kubefake.NewClientset(c.objects...)
client.PrependReactor("create", "serviceaccounts/token",
func(action clienttesting.Action) (handled bool, ret runtime.Object, err error) {
act, ok := action.(clienttesting.CreateActionImpl)
if !ok {
return false, nil, nil
}
tokenReq, ok := act.Object.(*authenticationv1.TokenRequest)
if !ok {
return false, nil, nil
}
tokenReq.Status.Token = "token"
return true, tokenReq, nil
},
)
config := &chart.KlusterletChartConfig{}
config, err := RenderBootstrapHubKubeConfig(client, c.apiserverURL)(context.TODO(), config)
if err != nil {
t.Fatalf("failed to render bootstrap hub kubeconfig: %v", err)
}
kConfig, err := clientcmd.NewClientConfigFromBytes([]byte(config.BootstrapHubKubeConfig))
if err != nil {
t.Fatalf("failed to load bootstrap hub kubeconfig: %v", err)
}
rawConfig, err := kConfig.RawConfig()
if err != nil {
t.Fatalf("failed to load bootstrap hub kubeconfig: %v", err)
}
cluster := rawConfig.Contexts[rawConfig.CurrentContext].Cluster
if rawConfig.Clusters[cluster].Server != c.expectedURL {
t.Errorf("apiserver is not rendered correctly")
}
})
}
}
+30
View File
@@ -31,6 +31,10 @@ import (
"open-cluster-management.io/ocm/pkg/registration/hub/clusterprofile"
"open-cluster-management.io/ocm/pkg/registration/hub/clusterrole"
"open-cluster-management.io/ocm/pkg/registration/hub/gc"
"open-cluster-management.io/ocm/pkg/registration/hub/importer"
importeroptions "open-cluster-management.io/ocm/pkg/registration/hub/importer/options"
cloudproviders "open-cluster-management.io/ocm/pkg/registration/hub/importer/providers"
"open-cluster-management.io/ocm/pkg/registration/hub/importer/providers/capi"
"open-cluster-management.io/ocm/pkg/registration/hub/lease"
"open-cluster-management.io/ocm/pkg/registration/hub/managedcluster"
"open-cluster-management.io/ocm/pkg/registration/hub/managedclusterset"
@@ -44,6 +48,7 @@ import (
type HubManagerOptions struct {
ClusterAutoApprovalUsers []string
GCResourceList []string
ImportOption *importeroptions.Options
}
// NewHubManagerOptions returns a HubManagerOptions
@@ -51,6 +56,7 @@ func NewHubManagerOptions() *HubManagerOptions {
return &HubManagerOptions{
GCResourceList: []string{"addon.open-cluster-management.io/v1alpha1/managedclusteraddons",
"work.open-cluster-management.io/v1/manifestworks"},
ImportOption: importeroptions.New(),
}
}
@@ -62,6 +68,7 @@ func (m *HubManagerOptions) AddFlags(fs *pflag.FlagSet) {
"A list GVR user can customize which are cleaned up after cluster is deleted. Format is group/version/resource, "+
"and the default are managedclusteraddon and manifestwork. The resources will be deleted in order."+
"The flag works only when ResourceCleanup feature gate is enable.")
m.ImportOption.AddFlags(fs)
}
// RunControllerManager starts the controllers on hub to manage spoke cluster registration.
@@ -244,6 +251,23 @@ func (m *HubManagerOptions) RunControllerManagerWithInformers(
)
}
var providers []cloudproviders.Interface
var clusterImporter factory.Controller
if features.HubMutableFeatureGate.Enabled(ocmfeature.ClusterImporter) {
providers = []cloudproviders.Interface{
capi.NewCAPIProvider(controllerContext.KubeConfig, clusterInformers.Cluster().V1().ManagedClusters()),
}
clusterImporter = importer.NewImporter(
[]importer.KlusterletConfigRenderer{
importer.RenderBootstrapHubKubeConfig(kubeClient, m.ImportOption.APIServerURL),
},
clusterClient,
clusterInformers.Cluster().V1().ManagedClusters(),
providers,
controllerContext.EventRecorder,
)
}
gcController := gc.NewGCController(
kubeInformers.Rbac().V1().ClusterRoles().Lister(),
kubeInformers.Rbac().V1().ClusterRoleBindings().Lister(),
@@ -284,6 +308,12 @@ func (m *HubManagerOptions) RunControllerManagerWithInformers(
if features.HubMutableFeatureGate.Enabled(ocmfeature.ClusterProfile) {
go clusterProfileController.Run(ctx, 1)
}
if features.HubMutableFeatureGate.Enabled(ocmfeature.ClusterImporter) {
for _, provider := range providers {
go provider.Run(ctx)
}
go clusterImporter.Run(ctx, 1)
}
go gcController.Run(ctx, 1)
+3
View File
@@ -71,6 +71,9 @@ var _ = ginkgo.BeforeSuite(func() {
// enable resourceCleanup feature gate
err = features.HubMutableFeatureGate.Set("ResourceCleanup=true")
gomega.Expect(err).NotTo(gomega.HaveOccurred())
err = features.HubMutableFeatureGate.Set("ClusterImporter=true")
gomega.Expect(err).NotTo(gomega.HaveOccurred())
})
var _ = ginkgo.AfterSuite(func() {
@@ -75,6 +75,8 @@ var CRDPaths = []string{
"./vendor/open-cluster-management.io/api/cluster/v1beta2/0000_01_clusters.open-cluster-management.io_managedclustersetbindings.crd.yaml",
// spoke
"./vendor/open-cluster-management.io/api/cluster/v1alpha1/0000_02_clusters.open-cluster-management.io_clusterclaims.crd.yaml",
// external API deps
"./test/integration/testdeps/capi/cluster.x-k8s.io_clusters.yaml",
}
func runAgent(name string, opt *spoke.SpokeAgentOptions, commOption *commonoptions.AgentOptions, cfg *rest.Config) context.CancelFunc {
@@ -196,12 +198,17 @@ var _ = ginkgo.BeforeSuite(func() {
err = features.HubMutableFeatureGate.Set("ResourceCleanup=true")
gomega.Expect(err).NotTo(gomega.HaveOccurred())
// enable clusterImporter feature gate
err = features.HubMutableFeatureGate.Set("ClusterImporter=true")
gomega.Expect(err).NotTo(gomega.HaveOccurred())
// start hub controller
var ctx context.Context
startHub = func() {
ctx, stopHub = context.WithCancel(context.Background())
go func() {
m := hub.NewHubManagerOptions()
m.ImportOption.APIServerURL = cfg.Host
m.ClusterAutoApprovalUsers = []string{util.AutoApprovalBootstrapUser}
err := m.RunControllerManager(ctx, &controllercmd.ControllerContext{
KubeConfig: cfg,
@@ -0,0 +1,183 @@
package registration_test
import (
"context"
"fmt"
"github.com/onsi/ginkgo/v2"
"github.com/onsi/gomega"
"github.com/openshift/api"
"github.com/openshift/library-go/pkg/operator/events"
"github.com/openshift/library-go/pkg/operator/resource/resourceapply"
corev1 "k8s.io/api/core/v1"
rbacv1 "k8s.io/api/rbac/v1"
apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/serializer"
"k8s.io/apimachinery/pkg/util/rand"
utilruntime "k8s.io/apimachinery/pkg/util/runtime"
"k8s.io/client-go/dynamic"
operatorclient "open-cluster-management.io/api/client/operator/clientset/versioned"
clusterv1 "open-cluster-management.io/api/cluster/v1"
operatorv1 "open-cluster-management.io/api/operator/v1"
testingcommon "open-cluster-management.io/ocm/pkg/common/testing"
"open-cluster-management.io/ocm/pkg/operator/helpers/chart"
"open-cluster-management.io/ocm/pkg/registration/hub/importer"
"open-cluster-management.io/ocm/pkg/registration/hub/importer/providers/capi"
"open-cluster-management.io/ocm/test/integration/util"
)
var (
genericScheme = runtime.NewScheme()
genericCodecs = serializer.NewCodecFactory(genericScheme)
genericCodec = genericCodecs.UniversalDeserializer()
)
func init() {
utilruntime.Must(api.InstallKube(genericScheme))
utilruntime.Must(apiextensionsv1.AddToScheme(genericScheme))
utilruntime.Must(operatorv1.Install(genericScheme))
}
var _ = ginkgo.Describe("Cluster Auto Importer", func() {
var managedClusterName string
var dynamicClient dynamic.Interface
var operatorClient operatorclient.Interface
ginkgo.BeforeEach(func() {
suffix := rand.String(5)
managedClusterName = fmt.Sprintf("managedcluster-%s", suffix)
ginkgo.By("Create bootstrap token")
clusterManagerConfig := chart.NewDefaultClusterManagerChartConfig()
clusterManagerConfig.CreateBootstrapSA = true
clusterManagerConfig.CreateNamespace = true
manifests, err := chart.RenderClusterManagerChart(clusterManagerConfig, "open-cluster-management")
gomega.Expect(err).NotTo(gomega.HaveOccurred())
recorder := events.NewInMemoryRecorder("importer-testing")
for _, manifest := range manifests {
requiredObj, _, err := genericCodec.Decode(manifest, nil, nil)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
switch t := requiredObj.(type) {
case *corev1.Namespace:
_, _, err = resourceapply.ApplyNamespace(context.TODO(), kubeClient.CoreV1(), recorder, t)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
case *corev1.ServiceAccount:
_, _, err = resourceapply.ApplyServiceAccount(context.TODO(), kubeClient.CoreV1(), recorder, t)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
case *rbacv1.ClusterRole:
_, _, err = resourceapply.ApplyClusterRole(context.TODO(), kubeClient.RbacV1(), recorder, t)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
case *rbacv1.ClusterRoleBinding:
_, _, err = resourceapply.ApplyClusterRoleBinding(context.TODO(), kubeClient.RbacV1(), recorder, t)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
}
}
dynamicClient, err = dynamic.NewForConfig(spokeCfg)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
operatorClient, err = operatorclient.NewForConfig(spokeCfg)
gomega.Expect(err).NotTo(gomega.HaveOccurred())
})
ginkgo.Context("Cluster API importer", func() {
ginkgo.JustBeforeEach(func() {
cluster := &clusterv1.ManagedCluster{
ObjectMeta: metav1.ObjectMeta{
Name: managedClusterName,
},
Spec: clusterv1.ManagedClusterSpec{
HubAcceptsClient: true,
},
}
_, err := clusterClient.ClusterV1().ManagedClusters().Create(context.TODO(), cluster, metav1.CreateOptions{})
gomega.Expect(err).NotTo(gomega.HaveOccurred())
})
ginkgo.JustAfterEach(func() {
err := clusterClient.ClusterV1().ManagedClusters().Delete(context.TODO(), managedClusterName, metav1.DeleteOptions{})
if !errors.IsNotFound(err) {
gomega.Expect(err).NotTo(gomega.HaveOccurred())
}
})
ginkgo.It("Should import CAPI cluster", func() {
ginkgo.By("Create CAPI cluster")
capiCluster := testingcommon.NewUnstructured(
"cluster.x-k8s.io/v1beta1", "Cluster", managedClusterName, managedClusterName)
_, err := dynamicClient.Resource(capi.ClusterAPIGVR).Namespace(managedClusterName).Create(
context.TODO(), capiCluster, metav1.CreateOptions{})
gomega.Expect(err).NotTo(gomega.HaveOccurred())
gomega.Eventually(func() error {
spokeCluster, err := util.GetManagedCluster(clusterClient, managedClusterName)
if err != nil {
return err
}
if !meta.IsStatusConditionFalse(
spokeCluster.Status.Conditions, importer.ManagedClusterConditionImported) {
return fmt.Errorf("cluster should have error when imported")
}
return nil
}, eventuallyTimeout, eventuallyInterval).ShouldNot(gomega.HaveOccurred())
ginkgo.By("Create secret")
capiSecret := &corev1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: managedClusterName + "-kubeconfig",
Namespace: managedClusterName,
},
Data: map[string][]byte{
"value": util.NewKubeConfig(spokeCfg),
},
}
_, err = kubeClient.CoreV1().Secrets(managedClusterName).Create(context.TODO(), capiSecret, metav1.CreateOptions{})
gomega.Expect(err).NotTo(gomega.HaveOccurred())
// trigger the capi cluster resource to reconcile again
gomega.Eventually(func() error {
capiCluster, err := dynamicClient.Resource(capi.ClusterAPIGVR).Namespace(managedClusterName).Get(
context.TODO(), managedClusterName, metav1.GetOptions{})
if err != nil {
return err
}
capiCluster.SetLabels(map[string]string{"reconcile": "trigger"})
_, err = dynamicClient.Resource(capi.ClusterAPIGVR).Namespace(managedClusterName).Update(
context.TODO(), capiCluster, metav1.UpdateOptions{})
if err != nil {
return err
}
return nil
}, eventuallyTimeout, eventuallyInterval).ShouldNot(gomega.HaveOccurred())
gomega.Eventually(func() error {
spokeCluster, err := util.GetManagedCluster(clusterClient, managedClusterName)
if err != nil {
return err
}
if !meta.IsStatusConditionTrue(
spokeCluster.Status.Conditions, importer.ManagedClusterConditionImported) {
return fmt.Errorf("cluster should have imported")
}
_, err = operatorClient.OperatorV1().Klusterlets().Get(
context.TODO(), "klusterlet", metav1.GetOptions{})
if err != nil {
return err
}
return nil
}, eventuallyTimeout, eventuallyInterval).ShouldNot(gomega.HaveOccurred())
err = dynamicClient.Resource(capi.ClusterAPIGVR).Namespace(managedClusterName).Delete(
context.TODO(), managedClusterName, metav1.DeleteOptions{})
gomega.Expect(err).NotTo(gomega.HaveOccurred())
err = kubeClient.CoreV1().Secrets(managedClusterName).Delete(
context.TODO(), managedClusterName+"-kubeconfig", metav1.DeleteOptions{})
gomega.Expect(err).NotTo(gomega.HaveOccurred())
})
})
})
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -1584,7 +1584,7 @@ open-cluster-management.io/addon-framework/pkg/agent
open-cluster-management.io/addon-framework/pkg/assets
open-cluster-management.io/addon-framework/pkg/index
open-cluster-management.io/addon-framework/pkg/utils
# open-cluster-management.io/api v0.15.1-0.20241209025232-b62746ae96d4
# open-cluster-management.io/api v0.15.1-0.20241210025410-0ba6809d0ae2
## explicit; go 1.22.0
open-cluster-management.io/api/addon/v1alpha1
open-cluster-management.io/api/client/addon/clientset/versioned
+6 -2
View File
@@ -80,6 +80,9 @@ const (
// ClusterProfile will start new controller in the Hub that can be used to sync ManagedCluster to ClusterProfile.
ClusterProfile featuregate.Feature = "ClusterProfile"
// ClusterImporter will enable the auto import of managed cluster for certain cluster providers, e.g. cluster-api.
ClusterImporter featuregate.Feature = "ClusterImporter"
)
// DefaultSpokeRegistrationFeatureGates consists of all known ocm-registration
@@ -87,7 +90,7 @@ const (
// add it here.
var DefaultSpokeRegistrationFeatureGates = map[featuregate.Feature]featuregate.FeatureSpec{
ClusterClaim: {Default: true, PreRelease: featuregate.Beta},
AddonManagement: {Default: true, PreRelease: featuregate.Alpha},
AddonManagement: {Default: true, PreRelease: featuregate.Beta},
V1beta1CSRAPICompatibility: {Default: false, PreRelease: featuregate.Alpha},
MultipleHubs: {Default: false, PreRelease: featuregate.Alpha},
}
@@ -101,10 +104,11 @@ var DefaultHubRegistrationFeatureGates = map[featuregate.Feature]featuregate.Fea
ManagedClusterAutoApproval: {Default: false, PreRelease: featuregate.Alpha},
ResourceCleanup: {Default: false, PreRelease: featuregate.Alpha},
ClusterProfile: {Default: false, PreRelease: featuregate.Alpha},
ClusterImporter: {Default: false, PreRelease: featuregate.Alpha},
}
var DefaultHubAddonManagerFeatureGates = map[featuregate.Feature]featuregate.FeatureSpec{
AddonManagement: {Default: true, PreRelease: featuregate.Alpha},
AddonManagement: {Default: true, PreRelease: featuregate.Beta},
}
// DefaultHubWorkFeatureGates consists of all known acm work wehbook feature keys.