Skip to content

Commit 4c62623

Browse files
committed
wip
1 parent 3d07014 commit 4c62623

2 files changed

Lines changed: 49 additions & 11 deletions

File tree

pkg/volume/csi/nodeinfomanager/nodeinfomanager.go

Lines changed: 32 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -74,7 +74,7 @@ type nodeUpdateFunc func(*v1.Node) (newNode *v1.Node, updated bool, err error)
7474

7575
// Interface implements an interface for managing labels of a node
7676
type Interface interface {
77-
CreateCSINode() (*storagev1.CSINode, error)
77+
CreateCSINode(node *v1.Node) (*storagev1.CSINode, error)
7878

7979
// Updates or Creates the CSINode object with annotations for CSI Migration
8080
InitializeCSINodeWithAnnotation() error
@@ -378,17 +378,36 @@ func (nim *nodeInfoManager) tryUpdateCSINode(
378378
maxAttachLimit int64,
379379
topology map[string]string) error {
380380

381+
node, err := csiKubeClient.CoreV1().Nodes().Get(context.TODO(), string(nim.nodeName), metav1.GetOptions{})
382+
if err != nil {
383+
return err
384+
}
385+
381386
nodeInfo, err := csiKubeClient.StorageV1().CSINodes().Get(context.TODO(), string(nim.nodeName), metav1.GetOptions{})
382387
if nodeInfo == nil || errors.IsNotFound(err) {
383-
nodeInfo, err = nim.CreateCSINode()
388+
nodeInfo, err = nim.CreateCSINode(node)
384389
}
390+
385391
if err != nil {
386392
return err
387393
}
388394

395+
if !nodeOwnsCSINode(node, nodeInfo) {
396+
return fmt.Errorf("CSINode %q is owned by different node", nodeInfo.Name)
397+
}
398+
389399
return nim.installDriverToCSINode(nodeInfo, driverName, driverNodeID, maxAttachLimit, topology)
390400
}
391401

402+
func nodeOwnsCSINode(node *v1.Node, nodeInfo *storagev1.CSINode) bool {
403+
for _, ownerRef := range nodeInfo.OwnerReferences {
404+
if ownerRef.Kind == nodeKind.Kind && ownerRef.Name == node.Name && ownerRef.UID == node.UID {
405+
return true
406+
}
407+
}
408+
return false
409+
}
410+
392411
func (nim *nodeInfoManager) InitializeCSINodeWithAnnotation() error {
393412
csiKubeClient := nim.volumeHost.GetKubeClient()
394413
if csiKubeClient == nil {
@@ -411,15 +430,24 @@ func (nim *nodeInfoManager) InitializeCSINodeWithAnnotation() error {
411430
}
412431

413432
func (nim *nodeInfoManager) tryInitializeCSINodeWithAnnotation(csiKubeClient clientset.Interface) error {
433+
node, err := csiKubeClient.CoreV1().Nodes().Get(context.TODO(), string(nim.nodeName), metav1.GetOptions{})
434+
if err != nil {
435+
return err
436+
}
437+
414438
nodeInfo, err := csiKubeClient.StorageV1().CSINodes().Get(context.TODO(), string(nim.nodeName), metav1.GetOptions{})
415439
if nodeInfo == nil || errors.IsNotFound(err) {
416440
// CreateCSINode will set the annotation
417-
_, err = nim.CreateCSINode()
441+
_, err = nim.CreateCSINode(node)
418442
return err
419443
} else if err != nil {
420444
return err
421445
}
422446

447+
if !nodeOwnsCSINode(node, nodeInfo) {
448+
return fmt.Errorf("CSINode %q is owned by different node", nodeInfo.Name)
449+
}
450+
423451
annotationModified := setMigrationAnnotation(nim.migratedPlugins, nodeInfo)
424452

425453
if annotationModified {
@@ -430,18 +458,13 @@ func (nim *nodeInfoManager) tryInitializeCSINodeWithAnnotation(csiKubeClient cli
430458

431459
}
432460

433-
func (nim *nodeInfoManager) CreateCSINode() (*storagev1.CSINode, error) {
461+
func (nim *nodeInfoManager) CreateCSINode(node *v1.Node) (*storagev1.CSINode, error) {
434462

435463
csiKubeClient := nim.volumeHost.GetKubeClient()
436464
if csiKubeClient == nil {
437465
return nil, fmt.Errorf("error getting CSI client")
438466
}
439467

440-
node, err := csiKubeClient.CoreV1().Nodes().Get(context.TODO(), string(nim.nodeName), metav1.GetOptions{})
441-
if err != nil {
442-
return nil, err
443-
}
444-
445468
nodeInfo := &storagev1.CSINode{
446469
ObjectMeta: metav1.ObjectMeta{
447470
Name: string(nim.nodeName),

pkg/volume/csi/nodeinfomanager/nodeinfomanager_test.go

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -307,6 +307,21 @@ func TestInstallCSIDriver(t *testing.T) {
307307
},
308308
},
309309
},
310+
{
311+
name: "pre-existing node info, but owned by previous node",
312+
existingNode: func() *v1.Node {
313+
node := generateNode(nil /*nodeIDs*/, nil /*labels*/, nil /*capacity*/)
314+
node.UID = types.UID("node1")
315+
return node
316+
}(),
317+
existingCSINode: func() *storage.CSINode {
318+
csiNode := generateCSINode(nil /*nodeIDs*/, nil /*volumeLimits*/, nil /*topologyKeys*/)
319+
csiNode.OwnerReferences[0].UID = types.UID("node2")
320+
return csiNode
321+
}(),
322+
inputNodeID: "com.example.csi/csi-node1",
323+
expectFail: true,
324+
},
310325
{
311326
name: "nil topology, empty node",
312327
driverName: "com.example.csi.driver1",
@@ -972,7 +987,7 @@ func TestInstallCSIDriverExistingAnnotation(t *testing.T) {
972987
nim := NewNodeInfoManager(types.NodeName(nodeName), host, nil)
973988

974989
// Act
975-
_, err = nim.CreateCSINode()
990+
_, err = nim.CreateCSINode(tc.existingNode)
976991
if err != nil {
977992
t.Errorf("expected no error from creating CSINodeinfo but got: %v", err)
978993
continue
@@ -1032,7 +1047,7 @@ func test(t *testing.T, addNodeInfo bool, testcases []testcase) {
10321047
nim := NewNodeInfoManager(types.NodeName(nodeName), host, nil)
10331048

10341049
//// Act
1035-
nim.CreateCSINode()
1050+
nim.CreateCSINode(tc.existingNode)
10361051
if addNodeInfo {
10371052
err = nim.InstallCSIDriver(tc.driverName, tc.inputNodeID, tc.inputVolumeLimit, tc.inputTopology)
10381053
} else {

0 commit comments

Comments
 (0)