Skip to content

Commit bdf648a

Browse files
authored
Merge pull request #3797 from landreasyan/cherry-pick-3668-release-1.33
[release-1.33] fix: use node informer to get node labels (cherry-pick of #3668)
2 parents bffac9f + 9feb380 commit bdf648a

14 files changed

Lines changed: 1187 additions & 20 deletions

File tree

pkg/azuredisk/azure_controller_common.go

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ import (
3131
"k8s.io/apimachinery/pkg/types"
3232
"k8s.io/apimachinery/pkg/util/wait"
3333
kwait "k8s.io/apimachinery/pkg/util/wait"
34+
"k8s.io/client-go/tools/cache"
3435
cloudprovider "k8s.io/cloud-provider"
3536
volerr "k8s.io/cloud-provider/volume/errors"
3637
"k8s.io/klog/v2"
@@ -88,6 +89,7 @@ type controllerCommon struct {
8889
lockMap *lockMap
8990
cloud *provider.Cloud
9091
clientFactory azclient.ClientFactory
92+
nodeLister cache.GenericLister
9193
// disk queue that is waiting for attach or detach on specific node
9294
// <nodeName, map<diskURI, *provider.AttachDiskOptions/DetachDiskOptions>>
9395
attachDiskMap sync.Map
@@ -224,7 +226,13 @@ func (c *controllerCommon) AttachDisk(ctx context.Context, diskName, diskURI str
224226

225227
numDisksAllowed := math.MaxInt
226228
if c.CheckDiskCountForBatching {
227-
_, instanceType, err := getNodeInfoFromLabels(ctx, string(nodeName), c.cloud.KubeClient)
229+
var instanceType string
230+
var err error
231+
if c.nodeLister != nil {
232+
_, instanceType, err = GetNodeInfoFromNodeLister(string(nodeName), c.nodeLister)
233+
} else {
234+
_, instanceType, err = GetNodeInfoFromLabels(ctx, string(nodeName), c.cloud.KubeClient)
235+
}
228236
if err != nil {
229237
klog.Errorf("failed to get node info from labels: %v", err)
230238
} else if instanceType != "" {

pkg/azuredisk/azuredisk.go

Lines changed: 89 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ import (
3636

3737
grpcprom "github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus"
3838
corev1 "k8s.io/api/core/v1"
39+
apierrors "k8s.io/apimachinery/pkg/api/errors"
3940
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
4041
"k8s.io/apimachinery/pkg/types"
4142
k8stypes "k8s.io/apimachinery/pkg/types"
@@ -50,6 +51,11 @@ import (
5051
"k8s.io/mount-utils"
5152
"k8s.io/utils/ptr"
5253

54+
"k8s.io/apimachinery/pkg/runtime/schema"
55+
"k8s.io/client-go/metadata"
56+
"k8s.io/client-go/metadata/metadatainformer"
57+
"k8s.io/client-go/tools/cache"
58+
5359
consts "sigs.k8s.io/azuredisk-csi-driver/pkg/azureconstants"
5460
"sigs.k8s.io/azuredisk-csi-driver/pkg/azureutils"
5561
csicommon "sigs.k8s.io/azuredisk-csi-driver/pkg/csi-common"
@@ -156,6 +162,9 @@ type Driver struct {
156162
enableMigrationMonitor bool
157163
// whether to convert ReadWrite cachingMode to ReadOnly for intree PVs to avoid issues
158164
convertRWCachingModeForIntreePV bool
165+
nodeLister cache.GenericLister
166+
nodeInformerSynced cache.InformerSynced
167+
nodeInformerFactory metadatainformer.SharedInformerFactory
159168
}
160169

161170
// NewDriver Creates a NewCSIDriver object. Assumes vendor version is equal to driver version &
@@ -240,10 +249,17 @@ func NewDriver(options *DriverOptions) *Driver {
240249
userAgent := GetUserAgent(driver.Name, driver.customUserAgent, driver.userAgentSuffix)
241250
klog.V(2).Infof("driver userAgent: %s", userAgent)
242251

243-
kubeClient, err := azureutils.GetKubeClient(options.Kubeconfig, options.KubeAPIQPS, options.KubeAPIBurst)
252+
kubeConfig, err := azureutils.GetKubeConfig(options.Kubeconfig, options.KubeAPIQPS, options.KubeAPIBurst)
244253
if err != nil {
245254
klog.Warningf("get kubeconfig(%s) failed with error: %v", options.Kubeconfig, err)
246255
}
256+
var kubeClient clientset.Interface
257+
if kubeConfig != nil {
258+
kubeClient, err = clientset.NewForConfig(kubeConfig)
259+
if err != nil {
260+
klog.Warningf("get kubeclient failed with error: %v", err)
261+
}
262+
}
247263
driver.kubeClient = kubeClient
248264

249265
cloud, err := azureutils.GetCloudProviderFromClient(context.Background(), kubeClient, driver.cloudConfigSecretName, driver.cloudConfigSecretNamespace,
@@ -325,6 +341,23 @@ func NewDriver(options *DriverOptions) *Driver {
325341
}
326342
}
327343

344+
if kubeConfig != nil && driver.checkDiskCountForBatching && driver.NodeID == "" {
345+
// Create a metadata-only node informer to cache node labels locally (controller only)
346+
metadataClient, err := metadata.NewForConfig(kubeConfig)
347+
if err != nil {
348+
klog.Warningf("failed to create metadata client: %v, node informer will not be used", err)
349+
} else {
350+
driver.nodeInformerFactory = metadatainformer.NewSharedInformerFactory(metadataClient, 0)
351+
nodeGVR := schema.GroupVersionResource{Group: "", Version: "v1", Resource: "nodes"}
352+
nodeInformer := driver.nodeInformerFactory.ForResource(nodeGVR)
353+
driver.nodeLister = nodeInformer.Lister()
354+
driver.nodeInformerSynced = nodeInformer.Informer().HasSynced
355+
if driver.diskController != nil {
356+
driver.diskController.nodeLister = driver.nodeLister
357+
}
358+
}
359+
}
360+
328361
driver.deviceHelper = optimization.NewSafeDeviceHelper()
329362

330363
if driver.getPerfOptimizationEnabled() {
@@ -413,6 +446,19 @@ func (d *Driver) Run(ctx context.Context) error {
413446
csi.RegisterControllerServer(s, d)
414447
csi.RegisterNodeServer(s, d)
415448

449+
// Start the node informer if it was set up during driver initialization
450+
if d.nodeInformerFactory != nil {
451+
d.nodeInformerFactory.Start(ctx.Done())
452+
syncCtx, syncCancel := context.WithTimeout(ctx, 30*time.Second)
453+
defer syncCancel()
454+
if !cache.WaitForCacheSync(syncCtx.Done(), d.nodeInformerSynced) {
455+
klog.Warningf("metadata node informer cache has not synced yet, will continue to sync in background")
456+
} else {
457+
klog.V(2).Infof("metadata node informer cache synced successfully")
458+
}
459+
klog.V(2).Infof("started metadata node informer for GetNodeInfoFromLabels caching")
460+
}
461+
416462
go func() {
417463
//graceful shutdown
418464
<-ctx.Done()
@@ -422,6 +468,11 @@ func (d *Driver) Run(ctx context.Context) error {
422468
d.migrationMonitor.Stop()
423469
}
424470

471+
// Shutdown node informer if it was started
472+
if d.nodeInformerFactory != nil {
473+
d.nodeInformerFactory.Shutdown()
474+
}
475+
425476
s.GracefulStop()
426477
}()
427478
// Driver d act as IdentityServer, ControllerServer and NodeServer
@@ -665,8 +716,38 @@ func (d *Driver) getUsedLunsFromNode(ctx context.Context, nodeName types.NodeNam
665716
return usedLuns, nil
666717
}
667718

668-
// getNodeInfoFromLabels get zone, instanceType from node labels
669-
func getNodeInfoFromLabels(ctx context.Context, nodeName string, kubeClient clientset.Interface) (string, string, error) {
719+
// GetNodeInfoFromNodeLister gets zone, instanceType from node labels using the cached nodeLister.
720+
func GetNodeInfoFromNodeLister(nodeName string, nodeLister cache.GenericLister) (string, string, error) {
721+
if nodeLister == nil {
722+
return "", "", fmt.Errorf("nodeLister is nil")
723+
}
724+
725+
obj, err := nodeLister.Get(nodeName)
726+
if err != nil {
727+
if apierrors.IsNotFound(err) {
728+
klog.V(4).Infof("GetNodeInfoFromNodeLister: node(%s) not found in lister cache", nodeName)
729+
return "", "", nil
730+
}
731+
return "", "", fmt.Errorf("get node(%s) from lister failed: %v", nodeName, err)
732+
}
733+
734+
pom, ok := obj.(*metav1.PartialObjectMetadata)
735+
if !ok {
736+
return "", "", fmt.Errorf("node(%s) from lister is not *metav1.PartialObjectMetadata", nodeName)
737+
}
738+
739+
if len(pom.Labels) == 0 {
740+
return "", "", fmt.Errorf("node(%s) label is empty", nodeName)
741+
}
742+
743+
zone := pom.Labels[consts.WellKnownTopologyKey]
744+
instanceType := pom.Labels[consts.InstanceTypeKey]
745+
klog.V(4).Infof("GetNodeInfoFromNodeLister: node(%s): zone=%s, instanceType=%s", nodeName, zone, instanceType)
746+
return zone, instanceType, nil
747+
}
748+
749+
// GetNodeInfoFromLabels gets zone, instanceType from node labels via the kubeClient API server.
750+
func GetNodeInfoFromLabels(ctx context.Context, nodeName string, kubeClient clientset.Interface) (string, string, error) {
670751
if kubeClient == nil || kubeClient.CoreV1() == nil {
671752
return "", "", fmt.Errorf("kubeClient is nil")
672753
}
@@ -679,7 +760,11 @@ func getNodeInfoFromLabels(ctx context.Context, nodeName string, kubeClient clie
679760
if len(node.Labels) == 0 {
680761
return "", "", fmt.Errorf("node(%s) label is empty", nodeName)
681762
}
682-
return node.Labels[consts.WellKnownTopologyKey], node.Labels[consts.InstanceTypeKey], nil
763+
764+
zone := node.Labels[consts.WellKnownTopologyKey]
765+
instanceType := node.Labels[consts.InstanceTypeKey]
766+
klog.V(4).Infof("GetNodeInfoFromLabels: node(%s) from API server: zone=%s, instanceType=%s", nodeName, zone, instanceType)
767+
return zone, instanceType, nil
683768
}
684769

685770
// getDefaultDiskIOPSReadWrite according to requestGiB

pkg/azuredisk/azuredisk_test.go

Lines changed: 151 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -30,8 +30,15 @@ import (
3030
"go.uber.org/mock/gomock"
3131
"golang.org/x/sync/errgroup"
3232
"google.golang.org/grpc/status"
33+
corev1 "k8s.io/api/core/v1"
34+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
35+
"k8s.io/apimachinery/pkg/labels"
36+
"k8s.io/apimachinery/pkg/runtime"
37+
"k8s.io/apimachinery/pkg/runtime/schema"
3338
"k8s.io/apimachinery/pkg/types"
3439
clientset "k8s.io/client-go/kubernetes"
40+
"k8s.io/client-go/kubernetes/fake"
41+
"k8s.io/client-go/tools/cache"
3542
"k8s.io/utils/ptr"
3643
consts "sigs.k8s.io/azuredisk-csi-driver/pkg/azureconstants"
3744
"sigs.k8s.io/cloud-provider-azure/pkg/azclient/diskclient/mock_diskclient"
@@ -290,24 +297,161 @@ func TestDriver_CheckDiskExists_Success(t *testing.T) {
290297
assert.Equal(t, err, nil)
291298
}
292299

300+
// fakeErrorNodeLister is a GenericLister that always returns a specified error.
301+
type fakeErrorNodeLister struct {
302+
err error
303+
}
304+
305+
func (f *fakeErrorNodeLister) List(selector labels.Selector) ([]runtime.Object, error) {
306+
return nil, f.err
307+
}
308+
309+
func (f *fakeErrorNodeLister) Get(name string) (runtime.Object, error) {
310+
return nil, f.err
311+
}
312+
313+
func (f *fakeErrorNodeLister) ByNamespace(namespace string) cache.GenericNamespaceLister {
314+
return nil
315+
}
316+
317+
func TestGetNodeInfoFromNodeLister(t *testing.T) {
318+
nodeGR := schema.GroupResource{Group: "", Resource: "nodes"}
319+
tests := []struct {
320+
name string
321+
nodeName string
322+
nodeLister cache.GenericLister
323+
expectedZone string
324+
expectedType string
325+
expectedError error
326+
expectErrorNil bool
327+
}{
328+
{
329+
name: "nil lister",
330+
nodeName: "node1",
331+
nodeLister: nil,
332+
expectedError: fmt.Errorf("nodeLister is nil"),
333+
},
334+
{
335+
name: "lister returns node with labels",
336+
nodeName: "node1",
337+
nodeLister: func() cache.GenericLister {
338+
indexer := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{})
339+
node := &metav1.PartialObjectMetadata{
340+
ObjectMeta: metav1.ObjectMeta{
341+
Name: "node1",
342+
Labels: map[string]string{
343+
consts.WellKnownTopologyKey: "westus2-1",
344+
consts.InstanceTypeKey: "Standard_DS2_v2",
345+
},
346+
},
347+
}
348+
_ = indexer.Add(node)
349+
return cache.NewGenericLister(indexer, nodeGR)
350+
}(),
351+
expectedZone: "westus2-1",
352+
expectedType: "Standard_DS2_v2",
353+
expectErrorNil: true,
354+
},
355+
{
356+
name: "lister returns node with empty labels",
357+
nodeName: "node1",
358+
nodeLister: func() cache.GenericLister {
359+
indexer := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{})
360+
node := &metav1.PartialObjectMetadata{
361+
ObjectMeta: metav1.ObjectMeta{
362+
Name: "node1",
363+
Labels: map[string]string{},
364+
},
365+
}
366+
_ = indexer.Add(node)
367+
return cache.NewGenericLister(indexer, nodeGR)
368+
}(),
369+
expectedError: fmt.Errorf("node(node1) label is empty"),
370+
},
371+
{
372+
name: "lister does not have the node - returns nil (NotFound is not an error)",
373+
nodeName: "missing-node",
374+
nodeLister: func() cache.GenericLister {
375+
indexer := cache.NewIndexer(cache.MetaNamespaceKeyFunc, cache.Indexers{})
376+
return cache.NewGenericLister(indexer, nodeGR)
377+
}(),
378+
expectedZone: "",
379+
expectedType: "",
380+
expectErrorNil: true,
381+
},
382+
{
383+
name: "lister returns non-NotFound error - propagates error",
384+
nodeName: "node1",
385+
nodeLister: &fakeErrorNodeLister{err: fmt.Errorf("connection refused")},
386+
expectedError: fmt.Errorf("get node(node1) from lister failed: connection refused"),
387+
},
388+
}
389+
390+
for _, test := range tests {
391+
t.Run(test.name, func(t *testing.T) {
392+
zone, instanceType, err := GetNodeInfoFromNodeLister(test.nodeName, test.nodeLister)
393+
if test.expectErrorNil {
394+
assert.NoError(t, err)
395+
assert.Equal(t, test.expectedZone, zone)
396+
assert.Equal(t, test.expectedType, instanceType)
397+
} else {
398+
assert.EqualError(t, err, test.expectedError.Error())
399+
}
400+
})
401+
}
402+
}
403+
293404
func TestGetNodeInfoFromLabels(t *testing.T) {
294405
tests := []struct {
295-
nodeName string
296-
kubeClient clientset.Interface
297-
expectedError error
406+
name string
407+
nodeName string
408+
kubeClient clientset.Interface
409+
expectedZone string
410+
expectedType string
411+
expectedError error
412+
expectErrorNil bool
298413
}{
299414
{
415+
name: "nil kubeClient",
300416
nodeName: "",
301417
kubeClient: nil,
302418
expectedError: fmt.Errorf("kubeClient is nil"),
303419
},
420+
{
421+
name: "kubeClient returns node with labels",
422+
nodeName: "node1",
423+
kubeClient: fake.NewSimpleClientset(&corev1.Node{
424+
ObjectMeta: metav1.ObjectMeta{
425+
Name: "node1",
426+
Labels: map[string]string{
427+
consts.WellKnownTopologyKey: "centralus-1",
428+
consts.InstanceTypeKey: "Standard_E8s_v3",
429+
},
430+
},
431+
}),
432+
expectedZone: "centralus-1",
433+
expectedType: "Standard_E8s_v3",
434+
expectErrorNil: true,
435+
},
436+
{
437+
name: "kubeClient node not found",
438+
nodeName: "missing-node",
439+
kubeClient: fake.NewSimpleClientset(),
440+
expectedError: fmt.Errorf("get node(missing-node) failed with nodes \"missing-node\" not found"),
441+
},
304442
}
305443

306444
for _, test := range tests {
307-
_, _, err := getNodeInfoFromLabels(context.TODO(), test.nodeName, test.kubeClient)
308-
if !reflect.DeepEqual(err, test.expectedError) {
309-
t.Errorf("Unexpected result: %v, expected result: %v", err, test.expectedError)
310-
}
445+
t.Run(test.name, func(t *testing.T) {
446+
zone, instanceType, err := GetNodeInfoFromLabels(context.TODO(), test.nodeName, test.kubeClient)
447+
if test.expectErrorNil {
448+
assert.NoError(t, err)
449+
assert.Equal(t, test.expectedZone, zone)
450+
assert.Equal(t, test.expectedType, instanceType)
451+
} else {
452+
assert.EqualError(t, err, test.expectedError.Error())
453+
}
454+
})
311455
}
312456
}
313457

0 commit comments

Comments
 (0)