Skip to content

Commit a514eb0

Browse files
authored
feat: expose a more basic member clusters discovery interface (#94)
* feat: expose a more basic member clusters discovery interface * fix lint
1 parent 25379d1 commit a514eb0

4 files changed

Lines changed: 121 additions & 137 deletions

File tree

multicluster/manager.go

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -27,8 +27,7 @@ import (
2727
"github.com/go-logr/logr"
2828
"k8s.io/apimachinery/pkg/runtime"
2929
"k8s.io/apimachinery/pkg/util/sets"
30-
"k8s.io/client-go/discovery/cached/memory"
31-
"k8s.io/client-go/kubernetes"
30+
"k8s.io/client-go/discovery"
3231
"k8s.io/client-go/kubernetes/scheme"
3332
"k8s.io/client-go/rest"
3433
"k8s.io/klog/v2/klogr"
@@ -162,11 +161,10 @@ func (m *Manager) addUpdateHandler(cluster string, cfg *rest.Config) (err error)
162161
m.mutex.RUnlock()
163162

164163
// Get MemCacheClient for the cluster
165-
clientset, err := kubernetes.NewForConfig(cfg)
164+
discoveryClient, err := discovery.NewDiscoveryClientForConfig(cfg)
166165
if err != nil {
167166
return err
168167
}
169-
clusterCachedDiscoveryClient := memory.NewMemCacheClient(clientset.Discovery())
170168

171169
// Create cache for the cluster
172170
mapper, err := apiutil.NewDynamicRESTMapper(cfg)
@@ -202,7 +200,7 @@ func (m *Manager) addUpdateHandler(cluster string, cfg *rest.Config) (err error)
202200

203201
m.log.Info("add cluster", "cluster", cluster)
204202
m.clusterCacheManager.AddClusterCache(cluster, clusterCache)
205-
m.clusterClientManager.AddClusterClient(cluster, delegatingClusterClient, clusterCachedDiscoveryClient)
203+
m.clusterClientManager.AddClusterClient(cluster, delegatingClusterClient, discoveryClient)
206204

207205
m.mutex.Lock()
208206
m.hasCluster[cluster] = struct{}{}

multicluster/manager_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -184,10 +184,10 @@ var _ = Describe("multicluster", func() {
184184
})
185185

186186
It("multiClusterClient get server groups and resources", func() {
187-
mcdiscovery, ok := clusterClient.(MultiClusterDiscovery)
187+
mcdiscovery, ok := clusterClient.(MultiClusterDiscoveryManager)
188188
Expect(ok).To(BeTrue())
189-
cachedDiscoveryClient := mcdiscovery.MembersCachedDiscoveryInterface()
190-
apiGroups, apiResourceLists, err := cachedDiscoveryClient.ServerGroupsAndResources()
189+
_, allDiscoveryClients := mcdiscovery.GetAllDiscoveryInterface()
190+
apiGroups, apiResourceLists, err := GetAllClusterServerGroupsAndResources(allDiscoveryClients)
191191
Expect(err).NotTo(HaveOccurred())
192192

193193
groupVersionSets := sets.NewString()

multicluster/multi_cluster_client.go

Lines changed: 5 additions & 129 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@ import (
2424

2525
"github.com/go-logr/logr"
2626
"k8s.io/apimachinery/pkg/api/meta"
27-
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
2827
"k8s.io/apimachinery/pkg/runtime"
2928
"k8s.io/apimachinery/pkg/types"
3029
"k8s.io/client-go/discovery"
@@ -37,28 +36,15 @@ import (
3736
"kusionstack.io/kube-utils/multicluster/metrics"
3837
)
3938

40-
// MultiClusterDiscovery provides fed and member clusters discovery interface
41-
type MultiClusterDiscovery interface {
42-
FedDiscoveryInterface() discovery.DiscoveryInterface
43-
MembersCachedDiscoveryInterface() PartialCachedDiscoveryInterface
44-
}
45-
46-
// PartialCachedDiscoveryInterface is a subset of discovery.DiscoveryInterface.
47-
type PartialCachedDiscoveryInterface interface {
48-
ServerGroupsAndResources() ([]*metav1.APIGroup, []*metav1.APIResourceList, error)
49-
Invalidate()
50-
Fresh() bool
51-
}
52-
5339
type ClusterClientManager interface {
54-
AddClusterClient(cluster string, clusterClient client.Client, clusterCachedDiscoveryClient discovery.CachedDiscoveryInterface)
40+
AddClusterClient(cluster string, clusterClient client.Client, discoveryClient discovery.DiscoveryInterface)
5541
RemoveClusterClient(cluster string)
5642
}
5743

5844
func MultiClusterClientBuilder(log logr.Logger) (cluster.NewClientFunc, ClusterClientManager) {
5945
mcc := &multiClusterClient{
6046
clusterToClient: map[string]client.Client{},
61-
clusterToDiscoveryClient: map[string]discovery.CachedDiscoveryInterface{},
47+
clusterToDiscoveryClient: map[string]discovery.DiscoveryInterface{},
6248

6349
log: log,
6450
}
@@ -97,7 +83,7 @@ func MultiClusterClientBuilder(log logr.Logger) (cluster.NewClientFunc, ClusterC
9783
var (
9884
_ client.Client = &multiClusterClient{}
9985

100-
_ MultiClusterDiscovery = &multiClusterClient{}
86+
_ MultiClusterDiscoveryManager = &multiClusterClient{}
10187

10288
_ ClusterClientManager = &multiClusterClient{}
10389
)
@@ -109,13 +95,13 @@ type multiClusterClient struct {
10995
fedMapper meta.RESTMapper
11096

11197
clusterToClient map[string]client.Client
112-
clusterToDiscoveryClient map[string]discovery.CachedDiscoveryInterface
98+
clusterToDiscoveryClient map[string]discovery.DiscoveryInterface
11399

114100
mutex sync.RWMutex
115101
log logr.Logger
116102
}
117103

118-
func (mcc *multiClusterClient) AddClusterClient(cluster string, clusterClient client.Client, clusterDiscoveryClient discovery.CachedDiscoveryInterface) {
104+
func (mcc *multiClusterClient) AddClusterClient(cluster string, clusterClient client.Client, clusterDiscoveryClient discovery.DiscoveryInterface) {
119105
mcc.mutex.Lock()
120106
defer mcc.mutex.Unlock()
121107

@@ -478,113 +464,3 @@ func (mcc *multiClusterClient) getClusterNames(ctx context.Context) (clusters []
478464
}
479465
return
480466
}
481-
482-
func (mcc *multiClusterClient) FedDiscoveryInterface() discovery.DiscoveryInterface {
483-
return mcc.fedDiscovery
484-
}
485-
486-
func (mcc *multiClusterClient) MembersCachedDiscoveryInterface() PartialCachedDiscoveryInterface {
487-
return &cachedMultiClusterDiscoveryClient{
488-
delegate: mcc,
489-
}
490-
}
491-
492-
type cachedMultiClusterDiscoveryClient struct {
493-
delegate *multiClusterClient
494-
}
495-
496-
// ServerGroupsAndResources returns the supported server groups and resources for all clusters.
497-
func (c *cachedMultiClusterDiscoveryClient) ServerGroupsAndResources() ([]*metav1.APIGroup, []*metav1.APIResourceList, error) {
498-
c.delegate.mutex.Lock()
499-
defer c.delegate.mutex.Unlock()
500-
501-
allDiscoveryClient := c.delegate.clusterToDiscoveryClient
502-
503-
// If there is only one cluster, we can use the cached discovery client to get the server groups and resources
504-
if len(allDiscoveryClient) == 1 {
505-
for _, cachedClient := range allDiscoveryClient {
506-
return cachedClient.ServerGroupsAndResources()
507-
}
508-
}
509-
510-
// If there are multiple clusters, we need to get the intersection of groups and resources
511-
var (
512-
groupVersionCount = make(map[string]int)
513-
groupVersionNameCount = make(map[string]int)
514-
515-
apiGroupsRes []*metav1.APIGroup
516-
apiResourceListsRes []*metav1.APIResourceList
517-
groupVersionToResources = make(map[string][]metav1.APIResource)
518-
)
519-
for _, cachedClient := range allDiscoveryClient {
520-
apiGroups, apiResourceLists, err := cachedClient.ServerGroupsAndResources()
521-
if err != nil {
522-
return nil, nil, err
523-
}
524-
525-
for _, apiGroup := range apiGroups {
526-
groupVersion := apiGroup.PreferredVersion.GroupVersion
527-
528-
if _, ok := groupVersionCount[groupVersion]; !ok {
529-
groupVersionCount[groupVersion] = 1
530-
} else {
531-
groupVersionCount[groupVersion]++
532-
533-
if groupVersionCount[groupVersion] == len(allDiscoveryClient) { // all clusters have this PreferredVersion
534-
apiGroupsRes = append(apiGroupsRes, apiGroup)
535-
}
536-
}
537-
}
538-
539-
for _, apiResourceList := range apiResourceLists {
540-
for i := range apiResourceList.APIResources {
541-
apiResource := apiResourceList.APIResources[i]
542-
groupVersionName := fmt.Sprintf("%s/%s", apiResourceList.GroupVersion, apiResource.Name)
543-
544-
if _, ok := groupVersionNameCount[groupVersionName]; !ok {
545-
groupVersionNameCount[groupVersionName] = 1
546-
} else {
547-
groupVersionNameCount[groupVersionName]++
548-
549-
if groupVersionNameCount[groupVersionName] == len(allDiscoveryClient) { // all clusters have this GroupVersion and Name
550-
groupVersionToResources[apiResourceList.GroupVersion] = append(groupVersionToResources[apiResourceList.GroupVersion], apiResource)
551-
}
552-
}
553-
}
554-
}
555-
}
556-
557-
for groupVersion, resources := range groupVersionToResources {
558-
apiResourceList := metav1.APIResourceList{
559-
TypeMeta: metav1.TypeMeta{Kind: "APIResourceList", APIVersion: "v1"},
560-
GroupVersion: groupVersion,
561-
}
562-
apiResourceList.APIResources = append(apiResourceList.APIResources, resources...)
563-
apiResourceListsRes = append(apiResourceListsRes, &apiResourceList)
564-
}
565-
566-
return apiGroupsRes, apiResourceListsRes, nil
567-
}
568-
569-
// Invalidate invalidates the cached discovery clients for all clusters.
570-
func (c *cachedMultiClusterDiscoveryClient) Invalidate() {
571-
c.delegate.mutex.Lock()
572-
defer c.delegate.mutex.Unlock()
573-
574-
for _, cachedClient := range c.delegate.clusterToDiscoveryClient {
575-
cachedClient.Invalidate()
576-
}
577-
}
578-
579-
// Fresh returns true if all cached discovery clients are fresh.
580-
func (c *cachedMultiClusterDiscoveryClient) Fresh() bool {
581-
c.delegate.mutex.Lock()
582-
defer c.delegate.mutex.Unlock()
583-
584-
for _, cachedClient := range c.delegate.clusterToDiscoveryClient {
585-
if !cachedClient.Fresh() {
586-
return false
587-
}
588-
}
589-
return true
590-
}
Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
/**
2+
* Copyright 2025 KusionStack Authors.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License");
5+
* you may not use this file except in compliance with the License.
6+
* You may obtain a copy of the License at
7+
*
8+
* http://www.apache.org/licenses/LICENSE-2.0
9+
*
10+
* Unless required by applicable law or agreed to in writing, software
11+
* distributed under the License is distributed on an "AS IS" BASIS,
12+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
* See the License for the specific language governing permissions and
14+
* limitations under the License.
15+
*/
16+
17+
package multicluster
18+
19+
import (
20+
"fmt"
21+
"maps"
22+
23+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
24+
"k8s.io/client-go/discovery"
25+
)
26+
27+
// MultiClusterDiscovery provides fed and member clusters discovery interface
28+
type MultiClusterDiscoveryManager interface {
29+
// GetAllDiscoveryInterface returns the fed and member clusters discovery interface
30+
GetAllDiscoveryInterface() (fed discovery.DiscoveryInterface, members map[string]discovery.DiscoveryInterface)
31+
}
32+
33+
func (mcc *multiClusterClient) GetAllDiscoveryInterface() (discovery.DiscoveryInterface, map[string]discovery.DiscoveryInterface) {
34+
mcc.mutex.RLock()
35+
defer mcc.mutex.RUnlock()
36+
members := make(map[string]discovery.DiscoveryInterface)
37+
maps.Copy(members, mcc.clusterToDiscoveryClient)
38+
return mcc.fedDiscovery, members
39+
}
40+
41+
func GetAllClusterServerGroupsAndResources(membersDiscoveryClients map[string]discovery.DiscoveryInterface) ([]*metav1.APIGroup, []*metav1.APIResourceList, error) {
42+
totalClusters := len(membersDiscoveryClients)
43+
if totalClusters == 0 {
44+
return nil, nil, nil
45+
}
46+
// If there is only one cluster, we can use the cached discovery client to get the server groups and resources
47+
if totalClusters == 1 {
48+
for _, cachedClient := range membersDiscoveryClients {
49+
return cachedClient.ServerGroupsAndResources()
50+
}
51+
}
52+
53+
// If there are multiple clusters, we need to get the intersection of groups and resources
54+
var (
55+
groupVersionCount = make(map[string]int)
56+
groupVersionNameCount = make(map[string]int)
57+
58+
apiGroupsRes []*metav1.APIGroup
59+
apiResourceListsRes []*metav1.APIResourceList
60+
groupVersionToResources = make(map[string][]metav1.APIResource)
61+
)
62+
for _, discoveryClient := range membersDiscoveryClients {
63+
apiGroups, apiResourceLists, err := discoveryClient.ServerGroupsAndResources()
64+
if err != nil {
65+
return nil, nil, err
66+
}
67+
68+
for _, apiGroup := range apiGroups {
69+
groupVersion := apiGroup.PreferredVersion.GroupVersion
70+
71+
if _, ok := groupVersionCount[groupVersion]; !ok {
72+
groupVersionCount[groupVersion] = 1
73+
} else {
74+
groupVersionCount[groupVersion]++
75+
76+
if groupVersionCount[groupVersion] == totalClusters { // all clusters have this PreferredVersion
77+
apiGroupsRes = append(apiGroupsRes, apiGroup)
78+
}
79+
}
80+
}
81+
82+
for _, apiResourceList := range apiResourceLists {
83+
for i := range apiResourceList.APIResources {
84+
apiResource := apiResourceList.APIResources[i]
85+
groupVersionName := fmt.Sprintf("%s/%s", apiResourceList.GroupVersion, apiResource.Name)
86+
87+
if _, ok := groupVersionNameCount[groupVersionName]; !ok {
88+
groupVersionNameCount[groupVersionName] = 1
89+
} else {
90+
groupVersionNameCount[groupVersionName]++
91+
92+
if groupVersionNameCount[groupVersionName] == totalClusters { // all clusters have this GroupVersion and Name
93+
groupVersionToResources[apiResourceList.GroupVersion] = append(groupVersionToResources[apiResourceList.GroupVersion], apiResource)
94+
}
95+
}
96+
}
97+
}
98+
}
99+
100+
for groupVersion, resources := range groupVersionToResources {
101+
apiResourceList := metav1.APIResourceList{
102+
TypeMeta: metav1.TypeMeta{Kind: "APIResourceList", APIVersion: "v1"},
103+
GroupVersion: groupVersion,
104+
}
105+
apiResourceList.APIResources = append(apiResourceList.APIResources, resources...)
106+
apiResourceListsRes = append(apiResourceListsRes, &apiResourceList)
107+
}
108+
109+
return apiGroupsRes, apiResourceListsRes, nil
110+
}

0 commit comments

Comments
 (0)