ext-provisioner2
# 1.provisioner
# 1.1.初始化
provisionController初始化复杂一些,会注册PVC/PV监听回调及各类workQueue,以监听及处理存储创建或删除请求。// creates a new provision controller using the given configuration parameters and with private informers. func NewProvisionController(...) *ProvisionController { ... // 实例化provisioner controller := &ProvisionController{ client: client, provisionerName: provisionerName, provisioner: provisioner, ... resyncPeriod: 15min, ... threadiness: 4, failedProvisionThreshold: 15, failedDeleteThreshold: 15, ... leaseDuration: 15s, renewDeadline: 10s, retryPeriod: 2s, ... } ... // workqueue controller.claimQueue = workqueue.NewNamedRateLimitingQueue(rateLimiter, "claims") controller.volumeQueue = workqueue.NewNamedRateLimitingQueue(rateLimiter, "volumes") // PersistentVolumeClaim claimHandler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { controller.enqueueClaim(obj) }, UpdateFunc: func(oldObj, newObj interface{}) { controller.enqueueClaim(newObj) }, ... } ... controller.claimInformer.AddEventHandlerWithResyncPeriod(claimHandler, controller.resyncPeriod) ... controller.claimInformer.AddIndexers(cache.Indexers{uidIndex: func(obj interface{}) ([]string, error) { uid, err := getObjectUID(obj) ... return []string{uid}, nil }}) ... controller.claimsIndexer = controller.claimInformer.GetIndexer() // PersistentVolume volumeHandler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { controller.enqueueVolume(obj) }, UpdateFunc: func(oldObj, newObj interface{}) { controller.enqueueVolume(newObj) }, DeleteFunc: func(obj interface{}) { controller.forgetVolume(obj) }, } ... controller.volumeInformer.AddEventHandlerWithResyncPeriod(volumeHandler, controller.resyncPeriod) ... controller.volumes = controller.volumeInformer.GetStore() ... // StorageClasses controller.classes = controller.classInformer.GetStore() // 异步限流 if controller.createProvisionerPVLimiter != nil { // 限速队列 controller.volumeStore = NewVolumeStoreQueue(client, createProvisionerPVLimiter, claimsIndexer, ...) // 同步阻塞 } else { if controller.createProvisionedPVBackoff == nil { ... // 初始化退避模块 controller.createProvisionedPVBackoff = &wait.Backoff{ Duration: 10s, Factor: 1, backoffSteps: 5 } } // 退避缓存 controller.volumeStore = NewBackoffStore(client, eventRecorder, createProvisionedPVBackoff, controller) } return controller }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
注意
provisionController会监听PVC/PV/StorageClass,基于事件监听触发存储创建或删除
# 1.2.run入口
provisionController.Run()会激活informer同步及监听,启动runClaimWorker和runVolumeWorker处理PVC/PV资源。// Run starts all of this controller's control loops func (ctrl *ProvisionController) Run(ctx context.Context) { run := func(ctx context.Context) { ... // 启动PVC/PV/SC同步 go ctrl.claimInformer.Run(ctx.Done()) go ctrl.volumeInformer.Run(ctx.Done()) go ctrl.classInformer.Run(ctx.Done()) // 同步完成检查 if !WaitForCacheSync(..., claimInformer.HasSynced, volumeInformer.HasSynced, classInformer.HasSynced) { return } // 启动4个worker for i := 0; i < ctrl.threadiness; i++ { // PVC处理 go wait.Until(func() { ctrl.runClaimWorker(ctx) }, time.Second, ctx.Done()) // PV处理 go wait.Until(func() { ctrl.runVolumeWorker(ctx) }, time.Second, ctx.Done()) } select {} } // 消费volumeStore go ctrl.volumeStore.Run(ctx, DefaultThreadiness) if ctrl.leaderElection { ... // 选举启动 leaderelection.RunOrDie(ctx, leaderelection.LeaderElectionConfig{ Lock: rl, LeaseDuration: 15s RenewDeadline: 10s RetryPeriod: 2s Callbacks: leaderelection.LeaderCallbacks{ OnStartedLeading: run, ... }, }) ... } else { run(ctx) } } func (q *queueStore) Run(ctx context.Context, threadiness int) { ... // 启动4个worker for i := 0; i < threadiness; i++ { go wait.Until(q.saveVolumeWorker, time.Second, ctx.Done()) } <-ctx.Done() }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
注意
provisionController.Run()的几个worker核心的是runClaimWorker和runVolumeWorker
# 1.3.saveVolume
volumeStore.saveVolumeWorker()消费待创建的PV队列项,调用q.doSaveVolume()执行PV创建及内部状态更新。func (q *queueStore) saveVolumeWorker() { for q.processNextWorkItem() { } } func (q *queueStore) processNextWorkItem() bool { obj, shutdown := q.queue.Get() defer q.queue.Done(obj) ... if volumeName, ok = obj.(string); !ok { q.queue.Forget(obj) return true } // 获取缓存的PV对象 volumeObj, found := q.volumes.Load(volumeName) if !found { q.queue.Forget(volumeName) return true } volume, ok := volumeObj.(*v1.PersistentVolume) if !ok { q.queue.Forget(volumeName) return true } // 创建PV if err := q.doSaveVolume(volume); err != nil { q.queue.AddRateLimited(volumeName) return true } q.volumes.Delete(volumeName) q.queue.Forget(volumeName) return true } func (q *queueStore) doSaveVolume(volume *v1.PersistentVolume) error { _, err := q.client.CoreV1().PersistentVolumes().Create(context.Background(), volume, metav1.CreateOptions{}) if err == nil || apierrs.IsAlreadyExists(err) { return nil } return fmt.Errorf("error saving volume %s: %s", volume.Name, err) }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
注意
volumeStore.saveVolumeWorker()负责的是PV创建,交互CSIDriver创建或删除存储由其它worker负责
# 2.claimWorker
# 2.1.runworker
ctrl.runClaimWorker()负责消费claimQueue,处理PVC对象的新增及更新,根据需要交互CSI创建存储及相关的PV对象。func (ctrl *ProvisionController) runClaimWorker(ctx context.Context) { for ctrl.processNextClaimWorkItem(ctx) { } } // processNextClaimWorkItem processes items from claimQueue func (ctrl *ProvisionController) processNextClaimWorkItem(ctx context.Context) bool { obj, shutdown := ctrl.claimQueue.Get() ... err := func() error { ... defer ctrl.claimQueue.Done(obj) ... if key, ok = obj.(string); !ok { ctrl.claimQueue.Forget(obj) return fmt.Errorf("expected string in workqueue but got %#v", obj) } // 执行创建 if err := ctrl.syncClaimHandler(ctx, key); err != nil { if ctrl.failedProvisionThreshold == 0 { ctrl.claimQueue.AddRateLimited(obj) } else if ctrl.claimQueue.NumRequeues(obj) < ctrl.failedProvisionThreshold { ctrl.claimQueue.AddRateLimited(obj) } else { ctrl.claimsInProgress.Delete(key) // leak a volume that's being provisioned in the background! } return fmt.Errorf("error syncing claim %q: %s", key, err.Error()) } ctrl.claimQueue.Forget(obj) // 更新进度调整 ctrl.claimsInProgress.Delete(key) return nil }() ... return true } // gets the claim from informer's cache then calls syncClaim. A non-nil error triggers requeuing of the claim. func (ctrl *ProvisionController) syncClaimHandler(ctx context.Context, key string) error { // 基于Indexer查询匹配的PV列表 objs, err := ctrl.claimsIndexer.ByIndex(uidIndex, key) ... if len(objs) > 0 { claimObj = objs[0] } else { // 由更新进度查找 obj, found := ctrl.claimsInProgress.Load(key) ... claimObj = obj } return ctrl.syncClaim(ctx, claimObj) } // checks if the claim should have a volume provisioned for it and provisions one if so. func (ctrl *ProvisionController) syncClaim(ctx context.Context, obj interface{}) error { claim, ok := obj.(*v1.PersistentVolumeClaim) ... // 检查volume创建条件 should, err := ctrl.shouldProvision(ctx, claim) ... if should { ... // 执行volume协调 status, err := ctrl.provisionClaimOperation(ctx, claim) ... if err == nil || status == ProvisioningFinished { ... // 清理进度 ctrl.claimsInProgress.Delete(string(claim.UID)) return err } // 正在处理,缓存进度 if status == ProvisioningInBackground { ctrl.claimsInProgress.Store(string(claim.UID), claim) } ... return err } return nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
注意
claimworker核心逻辑是ctrl.provisionClaimOperation(),负责基于创建或更新的PVC创建volume资源
# 2.2.provisionclaim
ctrl.provisionClaimOperation()交互provisioner创建相关的volume存储资源,调用volumeStore尝试创建关联的PV对象。// attempts to provision a volume for the given claim. func (ctrl *ProvisionController) provisionClaimOperation(...) (ProvisioningState, error) { // 获取关联的storageClass名称 claimClass := util.GetPersistentVolumeClaimClass(claim) ... // 生成PVName pvName := ctrl.getProvisionedVolumeNameForClaim(claim) // PV存在 _, exists, err := ctrl.volumes.GetByKey(pvName) if err == nil && exists { return ProvisioningFinished, errStopProvision } // 生成claimRef claimRef, err := ref.GetReference(scheme.Scheme, claim) ... // 块存储支持检查 if err = ctrl.canProvision(ctx, claim); err != nil { return ProvisioningFinished, errStopProvision } // 获取关联的storageClass class, err := ctrl.getStorageClass(claimClass) if err != nil { return ProvisioningFinished, err } // provisioner未知 if !ctrl.knownProvisioner(class.Provisioner) { return ProvisioningFinished, errStopProvision } ... // 关联Pod调度node if nodeName, ok := getString(claim.Annotations, annSelectedNode, annAlphaSelectedNode); ok { // 获取node if ctrl.nodeLister != nil { selectedNode, err = ctrl.nodeLister.Get(nodeName) } else { selectedNode, err = ctrl.client.CoreV1().Nodes().Get(ctx, nodeName, metav1.GetOptions{}) } if err != nil { // node不存在 if apierrs.IsNotFound(err) { // 清理PVC selectedNode及更新nodeInformer return ctrl.provisionVolumeErrorHandling(ctx, ProvisioningReschedule, err, claim, operation) } return ProvisioningNoChange, err } } ... // 创建volume及生成PV对象 volume, result, err := ctrl.provisioner.Provision(ctx, options) if err != nil { if ierr, ok := err.(*IgnoredError); ok { return ProvisioningFinished, errStopProvision } // 尝试清理PVC selectedNode及更新nodeInformer return ctrl.provisionVolumeErrorHandling(ctx, result, err, claim, operation) } // PV引用PVC volume.Spec.ClaimRef = claimRef // Add external provisioner finalizer if ctrl.addFinalizer && !ctrl.checkFinalizer(volume, finalizerPV) { volume.ObjectMeta.Finalizers = append(volume.ObjectMeta.Finalizers, finalizerPV) } // 设置provisioned-by注解 metav1.SetMetaDataAnnotation(&volume.ObjectMeta, annDynamicallyProvisioned, class.Provisioner) // 设置storageClass名称 volume.Spec.StorageClassName = claimClass // 尝试创建PV,失败则加入volumeStore缓存及队列触发重试 if err := ctrl.volumeStore.StoreVolume(claim, volume); err != nil { return ProvisioningFinished, err } // 创建成功,更新PV缓存 ctrl.volumes.Add(volume) ... return ProvisioningFinished, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
注意
ctrl.provisioner.Provision()负责交互CSI Driver接口创建volume存储
# 2.3.provision
ctrl.provisioner.Provision()会检查容量、存储类型,满足条件触发csi.CreateVolume()调用创建存储,返回关联的PV对象func (p *provisionWrapper) Provision(...) (pv PersistentVolume, state controller.ProvisioningState, err error) { // 执行创建 pv, state, err = p.Provisioner.Provision(ctx, options) if err == nil && pv != nil { if pv.Spec.NodeAffinity != nil { // sc对应的segment入队capacityQueue p.c.refreshTopology(*pv.Spec.NodeAffinity) } else if options.StorageClass != nil { // sc对应的segment入队capacityQueue p.c.refreshSC(options.StorageClass.Name) } } else if state != controller.ProvisioningNoChange { if options.StorageClass != nil { p.c.refreshSC(options.StorageClass.Name) } } return } func (p *csiProvisioner) Provision(...) (*v1.PersistentVolume, controller.ProvisioningState, error) { ... // 匹配provisioner provisioner, ok := claim.Annotations[annStorageProvisioner] if !ok { provisioner = claim.Annotations[annBetaStorageProvisioner] } if provisioner != p.driverName && claim.Annotations[annMigratedTo] != p.driverName { return nil, controller.ProvisioningFinished, &controller.IgnoredError{...} } // 绑定模式检查+拓扑兼容检查 owned, err := p.checkNode(ctx, claim, options.StorageClass, "provision") ... // 条件冲突 if !owned { return nil, controller.ProvisioningNoChange, &controller.IgnoredError{...} } // 检查及生成req result, state, err := p.prepareProvision(ctx, claim, options.StorageClass, options.SelectedNode) ... // 调用csiDriver执行volume创建 rep, err := p.csiClient.CreateVolume(createCtx, req) if err != nil { return nil, state, err } ... // volume属性 volumeAttributes := map[string]string{provisionerIDKey: p.identity} for k, v := range rep.Volume.VolumeContext { volumeAttributes[k] = v } // 容量 respCap := rep.GetVolume().GetCapacityBytes() // 部分CSIDriver不返回容量,用请求容量 if respCap == 0 { respCap = volSizeBytes // 创建容量低于请求容量 } else if respCap < volSizeBytes { ... // 重试5次清理volume cleanupVolume(ctx, p, delReq, provisionerCredentials) ... // use InBackground to retry the call, hoping the volume is deleted correctly next time. return nil, controller.ProvisioningInBackground, capErr } // clone/snopshot/跨NS引用 if options.PVC.Spec.DataSource != nil || (utilfeature.DefaultFeatureGate.Enabled(features.CrossNamespaceVolumeDataSource) && options.PVC.Spec.DataSourceRef != nil && options.PVC.Spec.DataSourceRef.Namespace != nil && len(*options.PVC.Spec.DataSourceRef.Namespace) > 0) { contentSource := rep.GetVolume().ContentSource // 源volume为空 if contentSource == nil { delReq := &csi.DeleteVolumeRequest{ VolumeId: rep.GetVolume().GetVolumeId(), } // 清理非法volume cleanupVolume(ctx, p, delReq, provisionerCredentials) ... return nil, controller.ProvisioningInBackground, sourceErr } } ... // 构造PV对象 pv := &v1.PersistentVolume{ ObjectMeta: metav1.ObjectMeta{ Name: pvName, }, Spec: v1.PersistentVolumeSpec{ AccessModes: options.PVC.Spec.AccessModes, MountOptions: options.StorageClass.MountOptions, Capacity: v1.ResourceList{ v1.ResourceName(v1.ResourceStorage): bytesToQuantity(respCap), }, // TODO wait for CSI VolumeSource API PersistentVolumeSource: v1.PersistentVolumeSource{ CSI: result.csiPVSource, }, }, } // 更新删除密钥注解 if result.provDeletionSecrets != nil { metav1.SetAnnotation(&pv.ObjectMeta, annDeletionSecretRefName, provDeletionSecrets.name) metav1.SetAnnotation(&pv.ObjectMeta, annDeletionSecretRefNamespace, provDeletionSecrets.namespace) } else { metav1.SetAnnotation(&pv.ObjectMeta, annDeletionSecretRefName, "") metav1.SetAnnotation(&pv.ObjectMeta, annDeletionSecretRefNamespace, "") } // 设置回收策略 if options.StorageClass.ReclaimPolicy != nil { pv.Spec.PersistentVolumeReclaimPolicy = *options.StorageClass.ReclaimPolicy } // 基于拓扑生成节点亲和 if p.supportsTopology() { pv.Spec.NodeAffinity = GenerateVolumeNodeAffinity(rep.Volume.AccessibleTopology) } // 设置卷模式 if options.PVC.Spec.VolumeMode != nil { pv.Spec.VolumeMode = options.PVC.Spec.VolumeMode } // 设置文件系统类型 if !util.CheckPersistentVolumeClaimModeBlock(options.PVC) { pv.Spec.PersistentVolumeSource.CSI.FSType = result.fsType } // 迁移卷 if result.migratedVolume { // 转换为in-tree格式 pv, err = p.translator.TranslateCSIPVToInTree(pv) // 转换失败 if err != nil { // 尝试清理volume err := p.Delete(ctx, pv) ... return nil, controller.ProvisioningFinished, err } } return pv, controller.ProvisioningFinished, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
注意
provisioner.Provision()本质上会进行一系列条件检查,组装createVolume请求参数调用csiClient.CreateVolume创建存储卷
# 2.4.prepare
p.prepareProvision()负责调用CreateVolume前校验PVC/StorageClass/Node,生成完整、合法及可执行的createVolumeRequest。// prepareProvision does non-destructive parameter checking and preparations for provisioning a volume. func (p *csiProvisioner) prepareProvision(...) (*prepareProvisionResult, controller.ProvisioningState, error) { // storageClass不能为空 if sc == nil { return nil, controller.ProvisioningFinished, errors.New("storage class was nil") } // 获取关联datasource dataSource, err := p.dataSource(ctx, claim) ... // 能力检查(controller-service/volume/snapshot/clone) p.checkDriverCapabilities(rc) ... // pvc selector不支持 if claim.Spec.Selector != nil { return nil, controller.ProvisioningFinished, fmt.Errorf("claim Selector is not supported") } // 生成volume名称 pvName, err := makeVolumeName(p.volumeNamePrefix, string(claim.ObjectMeta.UID), p.volumeNameUUIDLength) ... for k, v := range sc.Parameters { if strings.ToLower(k) == "fstype" || k == prefixedFsTypeKey { fsType = v fsTypesFound++ } ... } // 限制一种fsType if fsTypesFound > 1 { return nil, controller.ProvisioningFinished, fmt.Errorf(...) } // 默认fsType if fsType == "" && p.defaultFSType != "" { fsType = p.defaultFSType } // 请求容量 capacity := claim.Spec.Resources.Requests[v1.ResourceName(v1.ResourceStorage)] volSizeBytes := capacity.Value() // 生成volume能力描述 volumeCaps, err := p.getVolumeCapabilities(claim, sc, fsType) ... // 构造volume创建请求 req := csi.CreateVolumeRequest{ Name: pvName, Parameters: sc.Parameters, VolumeCapabilities: volumeCaps, CapacityRange: &csi.CapacityRange{ RequiredBytes: int64(volSizeBytes), }, } // clone/snapshot数据源 if dataSource != nil && (rc.clone || rc.snapshot) { // 由PVC/snapshot获取源volume volumeContentSource, err := p.getVolumeContentSource(ctx, claim, sc, dataSource) ... req.VolumeContentSource = volumeContentSource } // clone数据源 if dataSource != nil && rc.clone { // 设置cloning-protection finalizer p.setCloneFinalizer(ctx, claim, dataSource) ... } // CSI支持拓扑 if p.supportsTopology() { // volume必须及偏好满足的拓扑条件 requirements, err := GenerateAccessibilityRequirements(...) ... req.AccessibilityRequirements = requirements } // 解析provisionsecretRef provisionerSecretRef, err := getSecretReference(provisionerSecretParams, sc.Parameters, pvName, &v1.PersistentVolumeClaim{ ObjectMeta: metav1.ObjectMeta{ Name: claim.Name, Namespace: claim.Namespace } }) ... // 获取provision secret凭证 provisionerCredentials, err := getCredentials(ctx, p.client, provisionerSecretRef) ... req.Secrets = provisionerCredentials // 解析controllerPublish secretRef controllerPublishSecretRef, err := getSecretReference(publishSecretParams, sc.Parameters, pvName, claim) ... // 解析nodeState secretRef nodeStageSecretRef, err := getSecretReference(nodeStageSecretParams, sc.Parameters, pvName, claim) ... // 解析nodePublish secretRef nodePublishSecretRef, err := getSecretReference(nodePublishSecretParams, sc.Parameters, pvName, claim) ... // 解析controllerExpand secretRef controllerExpandSecretRef, err := getSecretReference(controllerExpandSecretParams, sc.Parameters, pvName, claim) ... // 解析nodeExpand secretRef nodeExpandSecretRef, err := getSecretReference(nodeExpandSecretParams, sc.Parameters, pvName, claim) ... // volume凭证信息对象 csiPVSource := &v1.CSIPersistentVolumeSource{ Driver: p.driverName, // VolumeHandle and VolumeAttributes will be added after provisioning. ControllerPublishSecretRef: controllerPublishSecretRef, NodeStageSecretRef: nodeStageSecretRef, NodePublishSecretRef: nodePublishSecretRef, ControllerExpandSecretRef: controllerExpandSecretRef, NodeExpandSecretRef: nodeExpandSecretRef, } // 清理特殊前缀参数,CSI不关心 req.Parameters, err = removePrefixedParameters(sc.Parameters) ... // 启用额外字段 if p.extraCreateMetadata { req.Parameters[pvcNameKey] = claim.GetName() req.Parameters[pvcNamespaceKey] = claim.GetNamespace() req.Parameters[pvNameKey] = pvName } ... // 缓存删除volume使用的secret if provisionerSecretRef != nil { deletionAnnSecrets.name = provisionerSecretRef.Name deletionAnnSecrets.namespace = provisionerSecretRef.Namespace } return &prepareProvisionResult{ fsType: fsType, migratedVolume: migratedVolume, req: &req, csiPVSource: csiPVSource, provDeletionSecrets: deletionAnnSecrets, }, controller.ProvisioningNoChange, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
注意
p.prepare()负责准备createVolume请求及相关的secret凭证
# 3.volumeWorker
# 3.1.runworker
volumeWorker负责处理volumeQueue,监听PV对象的ADD/Update事件,根据需要交互CSIDriver对账清理未使用的volume以避免泄漏。func (ctrl *ProvisionController) runVolumeWorker(ctx context.Context) { for ctrl.processNextVolumeWorkItem(ctx) { } } // processNextVolumeWorkItem processes items from volumeQueue func (ctrl *ProvisionController) processNextVolumeWorkItem(ctx context.Context) bool { obj, shutdown := ctrl.volumeQueue.Get() ... err := func() error { ... defer ctrl.volumeQueue.Done(obj) ... // item非法 if key, ok = obj.(string); !ok { ctrl.volumeQueue.Forget(obj) return fmt.Errorf("expected string in workqueue but got %#v", obj) } // volume同步 if err := ctrl.syncVolumeHandler(ctx, key); err != nil { if ctrl.failedDeleteThreshold == 0 { ctrl.volumeQueue.AddRateLimited(obj) } else if ctrl.volumeQueue.NumRequeues(obj) < ctrl.failedDeleteThreshold { ctrl.volumeQueue.AddRateLimited(obj) } ... return fmt.Errorf("error syncing volume %q: %s", key, err.Error()) } ctrl.volumeQueue.Forget(obj) return nil }() ... return true }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
注意
volume同步逻辑主要位于ctrl.syncVolumeHandler(),必要时清理未使用的volume
# 3.2.syncvolume
ctrl.syncVolumeHandler()会调用ctrl.shouldDelete()检查volume清理条件,执行ctrl.deleteVolumeOperation()进行回收。// syncVolumeHandler gets the volume from informer's cache then calls syncVolume func (ctrl *ProvisionController) syncVolumeHandler(ctx context.Context, key string) error { // 获取PV volumeObj, exists, err := ctrl.volumes.GetByKey(key) ... // 未找到 if !exists { return nil } return ctrl.syncVolume(ctx, volumeObj) } // syncVolume checks if the volume should be deleted and deletes if so func (ctrl *ProvisionController) syncVolume(ctx context.Context, obj interface{}) error { volume, ok := obj.(*v1.PersistentVolume) ... // 匹配PV对应provisioner if !ctrl.isProvisionerForVolume(ctx, volume) { // Current provisioner is not responsible for the volume return nil } // finalizer处理 volume, err := ctrl.handleProtectionFinalizer(ctx, volume) ... if ctrl.shouldDelete(ctx, volume) { return ctrl.deleteVolumeOperation(ctx, volume) } return nil } func (ctrl *ProvisionController) handleProtectionFinalizer(...) (*v1.PersistentVolume, error) { // 回收策略及finalizers reclaimPolicy := volume.Spec.PersistentVolumeReclaimPolicy volumeFinalizers := volume.ObjectMeta.Finalizers // 启用PV回收限制+删除回收策略+未删除+Bound状态 if ctrl.addFinalizer && reclaimPolicy == v1.PersistentVolumeReclaimDelete && volume.DeletionTimestamp == nil && volume.Status.Phase == v1.VolumeBound { // 补充external-provisioner finalizer volumeFinalizers, modified = addFinalizer(volumeFinalizers, finalizerPV) } // 未启用删除限制 || 保留策略 || 复用策略 if !ctrl.addFinalizer || reclaimPolicy == v1.PersistentVolumeReclaimRetain || reclaimPolicy == v1.PersistentVolumeReclaimRecycle { // 清理external-provisioner finalizer volumeFinalizers, modified = removeFinalizer(volumeFinalizers, finalizerPV) } // PV调整 if modified { // 更新finalizer volume.ObjectMeta.Finalizers = volumeFinalizers newVolume, err := ctrl.updatePersistentVolume(ctx, volume) ... volume = newVolume } return volume, nil }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
注意
ctrl.syncVolumeHandler()会更新PV finalizer,检查删除条件及回收volume
# 3.3.checkdelete
ctrl.shouldDelete()检查provisioner删除授权,判断PV回收策略,匹配driver provisioner,校验volume的合法删除条件。// returns whether a volume should have its backing volume deleted. func (ctrl *ProvisionController) shouldDelete(ctx context.Context, volume *v1.PersistentVolume) bool { // 删除授权检查(未实现) if deletionGuard, ok := ctrl.provisioner.(DeletionGuard); ok { if !deletionGuard.ShouldDelete(ctx, volume) { return false } } // 启用PV删除限制 if ctrl.addFinalizer { // 已经处理过删除 if !ctrl.checkFinalizer(volume, finalizerPV) && volume.ObjectMeta.DeletionTimestamp != nil { return false } // 未启用PV删除限制 } else { // PV正在删除,无法接管 if volume.ObjectMeta.DeletionTimestamp != nil { return false } } // 非released状态 if volume.Status.Phase != v1.VolumeReleased { return false } // 不是删除回收策略 if volume.Spec.PersistentVolumeReclaimPolicy != v1.PersistentVolumeReclaimDelete { return false } return true }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35注意
shouldDelete检查的主要是fianlizer限制及PV的回收策略及状态
# 3.4.deleteVolume
ctrl.deleteVolumeOperation()负责调用provisioner.Delete()删除volume,清理PV对象及摘除protection finalizer。// attempts to delete the volume backing the given volume. func (ctrl *ProvisionController) deleteVolumeOperation(ctx context.Context, volume *v1.PersistentVolume) error { // 执行volume删除 err := ctrl.provisioner.Delete(ctx, volume) if err != nil { if ierr, ok := err.(*IgnoredError); ok { // Delete ignored, do nothing and hope another provisioner will delete it. return nil } return err } // 删除PV对象 ctrl.client.CoreV1().PersistentVolumes().Delete(ctx, volume.Name, metav1.DeleteOptions{}) ... // 启用PV删除限制 if ctrl.addFinalizer { if len(volume.ObjectMeta.Finalizers) > 0 { // 获取volumes缓存的PV对象 volumeObj, exists, err := ctrl.volumes.GetByKey(volume.Name) ... // 未找到 if !exists { // If the volume is not found return return nil } // 清理external-provisioner finalizer newVolume, ok := volumeObj.(*v1.PersistentVolume) ... finalizers, modified := removeFinalizer(newVolume.ObjectMeta.Finalizers, finalizerPV) // 更新PV对象 if modified { newVolume.ObjectMeta.Finalizers = finalizers ctrl.client.CoreV1().PersistentVolumes().Update(ctx, newVolume, metav1.UpdateOptions{}) ... } } } return nil } func (p *provisionWrapper) Delete(ctx context.Context, pv *v1.PersistentVolume) (err error) { err = p.Provisioner.Delete(ctx, pv) if err == nil && pv.Spec.NodeAffinity != nil { // segment入队capacityQueue p.c.refreshTopology(*pv.Spec.NodeAffinity) } return } func (p *csiProvisioner) Delete(ctx context.Context, volume *v1.PersistentVolume) error { if volume == nil { return fmt.Errorf("invalid CSI PV") } ... // 迁移卷 if p.translator.IsPVMigratable(volume) { migratedVolume = true // 转为CSI格式PV volume, err = p.translator.TranslateInTreePVToCSI(volume) ... } if volume.Spec.CSI == nil { return fmt.Errorf("invalid CSI PV") } // DS部署模式 if p.nodeDeployment != nil { // PV是当前节点的可访问卷 accessible, err := VolumeIsAccessible(volume.Spec.NodeAffinity, p.nodeDeployment.NodeInfo.AccessibleTopology) ... if !accessible { return &controller.IgnoredError{ Reason: "PV was not provisioned on this node", } } } // 解析volumeID volumeId := p.volumeHandleToId(volume.Spec.CSI.VolumeHandle) ... // 能力检查(provision/attach/clone) p.checkDriverCapabilities(rc) ... // 注入deletion secret p.handleSecretsForDeletion(ctx, volume, &req, migratedVolume) ... // 确认无volumeAttachment绑定 p.canDeleteVolume(volume) ... // 执行删除 _, err = p.csiClient.DeleteVolume(deleteCtx, &req) return err }1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107注意
deleteVolume处理的是PVC删除触发的PV级联回收,若直接删除PV对象这里不会处理