diff --git a/go.mod b/go.mod index 660de639a..c674ee4ca 100644 --- a/go.mod +++ b/go.mod @@ -82,6 +82,7 @@ require ( k8s.io/kube-aggregator v0.22.1 k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e k8s.io/kubectl v0.21.0 + k8s.io/metrics v0.21.0 k8s.io/utils v0.0.0-20210802155522-efc7438f0176 open-cluster-management.io/api v0.0.0-20210804091127-340467ff6239 rsc.io/letsencrypt v0.0.3 // indirect diff --git a/go.sum b/go.sum index 71986f652..397de5367 100644 --- a/go.sum +++ b/go.sum @@ -2589,6 +2589,7 @@ k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e h1:KLHHjkdQFomZy8+06csTWZ k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e/go.mod h1:vHXdDvt9+2spS2Rx9ql3I8tycm3H9FDfdUoIuKCefvw= k8s.io/kubectl v0.21.0 h1:WZXlnG/yjcE4LWO2g6ULjFxtzK6H1TKzsfaBFuVIhNg= k8s.io/kubectl v0.21.0/go.mod h1:EU37NukZRXn1TpAkMUoy8Z/B2u6wjHDS4aInsDzVvks= +k8s.io/metrics v0.21.0 h1:uwS3CgheLKaw3PTpwhjMswnm/PMqeLbdLH88VI7FMQQ= k8s.io/metrics v0.21.0/go.mod h1:L3Ji9EGPP1YBbfm9sPfEXSpnj8i24bfQbAFAsW0NueQ= k8s.io/utils v0.0.0-20190221042446-c2654d5206da/go.mod h1:8k8uAuAQ0rXslZKaEWd0c3oVhZz7sSzSiPnVZayjIX0= k8s.io/utils v0.0.0-20190801114015-581e00157fb1/go.mod h1:sZAwmy6armz5eXlNoLmJcl4F1QuKu7sr+mFQ0byX7Ew= diff --git a/pkg/multicluster/cluster_metrics_management.go b/pkg/multicluster/cluster_metrics_management.go new file mode 100644 index 000000000..33085ad59 --- /dev/null +++ b/pkg/multicluster/cluster_metrics_management.go @@ -0,0 +1,74 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package multicluster + +import ( + "context" + + "k8s.io/klog/v2" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// metricsMap records the metrics of clusters +var metricsMap map[string]*ClusterMetrics + +// ClusterMetricsMgr manage metrics of clusters +type ClusterMetricsMgr struct { + kubeClient client.Client +} + +// ClusterMetricsHelper is the interface that provides operations for cluster metrics +type ClusterMetricsHelper interface { + Refresh() error +} + +// NewClusterMetricsMgr will create a cluster metrics manager +func NewClusterMetricsMgr(kubeClient client.Client) (*ClusterMetricsMgr, error) { + mgr := &ClusterMetricsMgr{ + kubeClient: kubeClient, + } + err := mgr.Refresh() + return mgr, err +} + +// Refresh will re-collect cluster metrics and refresh cache +func (cmm *ClusterMetricsMgr) Refresh() error { + clusters, _ := ListVirtualClusters(context.Background(), cmm.kubeClient) + m := make(map[string]*ClusterMetrics) + + // retrieves metrics by cluster-gateway + for _, cluster := range clusters { + isConnected := true + clusterInfo, err := GetClusterInfo(context.Background(), cmm.kubeClient, cluster.Name) + if err != nil { + klog.Warningf("failed to get cluster info of cluster-(%s)", cluster.Name) + isConnected = false + } + clusterUsageMetrics, err := GetClusterMetricsFromMetricsAPI(context.Background(), cmm.kubeClient, cluster.Name) + if err != nil { + klog.Warningf("failed to request metrics api of cluster-(%s)", cluster.Name) + } + metrics := &ClusterMetrics{ + IsConnected: isConnected, + ClusterInfo: clusterInfo, + ClusterUsageMetrics: clusterUsageMetrics, + } + m[cluster.Name] = metrics + } + metricsMap = m + return nil +} diff --git a/pkg/multicluster/cluster_metrics_management_test.go b/pkg/multicluster/cluster_metrics_management_test.go new file mode 100644 index 000000000..b2e8ec4ea --- /dev/null +++ b/pkg/multicluster/cluster_metrics_management_test.go @@ -0,0 +1,155 @@ +/* +Copyright 2022 The KubeVela Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package multicluster + +import ( + "context" + "errors" + "strconv" + "testing" + + "gotest.tools/assert" + + "github.com/oam-dev/cluster-gateway/pkg/apis/cluster/v1alpha1" + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" + metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1" + clusterv1 "open-cluster-management.io/api/cluster/v1" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/fake" + + "github.com/oam-dev/kubevela/pkg/utils/common" +) + +const ( + NormalClusterName = "normal-cluster" + DisconnectedClusterName = "disconnected-cluster" + + NodeName1 = "node-1" + NodeName2 = "node-2" +) + +func TestRefresh(t *testing.T) { + fakeClient := NewFakeClient(fake.NewClientBuilder(). + WithScheme(common.Scheme). + WithRuntimeObjects(FakeManagedCluster("managed-cluster")). + WithObjects(FakeSecret(NormalClusterName), FakeSecret(DisconnectedClusterName)). + Build()) + + normalCluster := fake.NewClientBuilder(). + WithScheme(common.Scheme). + WithObjects(FakeNode(NodeName1, "8", strconv.FormatInt(16*1024*1024*1024, 10)), + FakeNode(NodeName2, "7", strconv.FormatInt(32*1024*1024*1024, 10)), + FakeNodeMetrics(NodeName1, "4", strconv.FormatInt(8*1024*1024*1024, 10)), + FakeNodeMetrics(NodeName2, "1", strconv.FormatInt(3*1024*1024*1024, 10))). + Build() + + disconnectedCluster := &disconnectedClient{} + + fakeClient.AddCluster(NormalClusterName, normalCluster) + fakeClient.AddCluster(DisconnectedClusterName, disconnectedCluster) + + mgr, err := NewClusterMetricsMgr(fakeClient) + assert.NilError(t, err) + + err = mgr.Refresh() + assert.NilError(t, err) + + clusters, err := ListVirtualClusters(context.Background(), fakeClient) + assert.NilError(t, err) + + for _, cluster := range clusters { + assertClusterMetrics(t, &cluster) + } + + disCluster, err := GetVirtualCluster(context.Background(), fakeClient, DisconnectedClusterName) + assert.NilError(t, err) + assertClusterMetrics(t, disCluster) + + norCluster, err := GetVirtualCluster(context.Background(), fakeClient, NormalClusterName) + assert.NilError(t, err) + assertClusterMetrics(t, norCluster) +} + +func assertClusterMetrics(t *testing.T, cluster *VirtualCluster) { + metrics := cluster.Metrics + switch cluster.Name { + case DisconnectedClusterName: + assert.Equal(t, metrics.IsConnected, false) + assert.Assert(t, metrics.ClusterInfo == nil) + assert.Assert(t, metrics.ClusterUsageMetrics == nil) + case NormalClusterName: + assert.Equal(t, metrics.IsConnected, true) + + assert.Assert(t, resource.MustParse("15").Equal(metrics.ClusterInfo.CPUCapacity)) + assert.Assert(t, resource.MustParse(strconv.FormatInt(48*1024*1024*1024, 10)).Equal(metrics.ClusterInfo.MemoryCapacity)) + assert.Assert(t, resource.MustParse("15").Equal(metrics.ClusterInfo.CPUAllocatable)) + assert.Assert(t, resource.MustParse(strconv.FormatInt(48*1024*1024*1024, 10)).Equal(metrics.ClusterInfo.MemoryAllocatable)) + + assert.Assert(t, resource.MustParse("5").Equal(metrics.ClusterUsageMetrics.CPUUsage)) + assert.Assert(t, resource.MustParse(strconv.FormatInt(11*1024*1024*1024, 10)).Equal(metrics.ClusterUsageMetrics.MemoryUsage)) + } +} + +func FakeNodeMetrics(name string, cpu string, memory string) *metricsV1beta1api.NodeMetrics { + nodeMetrics := &metricsV1beta1api.NodeMetrics{} + nodeMetrics.Name = name + nodeMetrics.Usage = corev1.ResourceList{ + corev1.ResourceCPU: resource.MustParse(cpu), + corev1.ResourceMemory: resource.MustParse(memory), + } + return nodeMetrics +} + +func FakeNode(name string, cpu string, memory string) *corev1.Node { + node := &corev1.Node{} + node.Name = name + node.Status = corev1.NodeStatus{ + Allocatable: map[corev1.ResourceName]resource.Quantity{ + corev1.ResourceCPU: resource.MustParse(cpu), + corev1.ResourceMemory: resource.MustParse(memory), + }, + Capacity: map[corev1.ResourceName]resource.Quantity{ + corev1.ResourceCPU: resource.MustParse(cpu), + corev1.ResourceMemory: resource.MustParse(memory), + }, + } + return node +} + +func FakeSecret(name string) *corev1.Secret { + secret := &corev1.Secret{} + secret.Name = name + secret.Labels = map[string]string{ + v1alpha1.LabelKeyClusterCredentialType: "ServiceAccountToken", + } + return secret +} + +func FakeManagedCluster(name string) *clusterv1.ManagedCluster { + managedCluster := &clusterv1.ManagedCluster{} + managedCluster.Name = name + return managedCluster +} + +type disconnectedClient struct { + client.Client +} + +func (cli *disconnectedClient) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error { + return errors.New("no such host") +} diff --git a/pkg/multicluster/o11n.go b/pkg/multicluster/o11n.go index 61b8da0a7..582504bc0 100644 --- a/pkg/multicluster/o11n.go +++ b/pkg/multicluster/o11n.go @@ -19,6 +19,8 @@ package multicluster import ( "context" + metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1" + "github.com/pkg/errors" corev1 "k8s.io/api/core/v1" storagev1 "k8s.io/api/storage/v1" @@ -26,6 +28,13 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client" ) +// ClusterMetrics describes the metrics of a cluster +type ClusterMetrics struct { + IsConnected bool + ClusterInfo *ClusterInfo + ClusterUsageMetrics *ClusterUsageMetrics +} + // ClusterInfo describes the basic information of a cluster type ClusterInfo struct { Nodes *corev1.NodeList @@ -40,6 +49,12 @@ type ClusterInfo struct { StorageClasses *storagev1.StorageClassList } +// ClusterUsageMetrics describes the usage metrics of a cluster +type ClusterUsageMetrics struct { + CPUUsage resource.Quantity + MemoryUsage resource.Quantity +} + // GetClusterInfo retrieves current cluster info from cluster func GetClusterInfo(_ctx context.Context, k8sClient client.Client, clusterName string) (*ClusterInfo, error) { ctx := ContextWithClusterName(_ctx, clusterName) @@ -81,3 +96,20 @@ func GetClusterInfo(_ctx context.Context, k8sClient client.Client, clusterName s StorageClasses: storageClasses, }, nil } + +// GetClusterMetricsFromMetricsAPI retrieves current cluster metrics based on GetNodeMetricsFromMetricsAPI +func GetClusterMetricsFromMetricsAPI(ctx context.Context, k8sClient client.Client, clusterName string) (*ClusterUsageMetrics, error) { + nodeMetricsList := metricsV1beta1api.NodeMetricsList{} + if err := k8sClient.List(ContextWithClusterName(ctx, clusterName), &nodeMetricsList); err != nil { + return nil, errors.Wrapf(err, "failed to list node metrics") + } + var memoryUsage, cpuUsage resource.Quantity + for _, nm := range nodeMetricsList.Items { + cpuUsage.Add(*nm.Usage.Cpu()) + memoryUsage.Add(*nm.Usage.Memory()) + } + return &ClusterUsageMetrics{ + CPUUsage: cpuUsage, + MemoryUsage: memoryUsage, + }, nil +} diff --git a/pkg/multicluster/virtual_cluster.go b/pkg/multicluster/virtual_cluster.go index dfea74142..a41adaf53 100644 --- a/pkg/multicluster/virtual_cluster.go +++ b/pkg/multicluster/virtual_cluster.go @@ -46,6 +46,7 @@ type VirtualCluster struct { EndPoint string Accepted bool Labels map[string]string + Metrics *ClusterMetrics } // NewVirtualClusterFromSecret extract virtual cluster from cluster secret @@ -68,6 +69,7 @@ func NewVirtualClusterFromSecret(secret *corev1.Secret) (*VirtualCluster, error) EndPoint: endpoint, Accepted: true, Labels: labels, + Metrics: metricsMap[secret.Name], }, nil } @@ -82,6 +84,7 @@ func NewVirtualClusterFromManagedCluster(managedCluster *clusterv1.ManagedCluste EndPoint: "-", Accepted: managedCluster.Spec.HubAcceptsClient, Labels: managedCluster.GetLabels(), + Metrics: metricsMap[managedCluster.Name], }, nil } diff --git a/pkg/utils/common/common.go b/pkg/utils/common/common.go index b444efdea..400684192 100644 --- a/pkg/utils/common/common.go +++ b/pkg/utils/common/common.go @@ -56,6 +56,8 @@ import ( "sigs.k8s.io/controller-runtime/pkg/client/config" "sigs.k8s.io/yaml" + metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1" + oamcore "github.com/oam-dev/kubevela/apis/core.oam.dev" "github.com/oam-dev/kubevela/apis/core.oam.dev/common" "github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1" @@ -96,6 +98,7 @@ func init() { _ = ocmclusterv1.Install(Scheme) _ = ocmworkv1.Install(Scheme) _ = clustergatewayapi.AddToScheme(Scheme) + _ = metricsV1beta1api.AddToScheme(Scheme) // +kubebuilder:scaffold:scheme }