ext-provisioner1
# 1.简介
# 1.1.作用
external-provisioner参与存储资源及PV对象的创建,负责监听PVC资源的创建事件,交互Custom CSI创建或回收存储及相关PV对象。
注意
rdb和cephfs共用一套处理逻辑,核心实现位于provisionController
# 1.2.入口
main函数作为入口会初始化capacityController、csiClaimController和provisionController关键组件,启动异步任务维护存储状态。func main() { ... // 建立grpc连接 grpcClient, err := ctrl.Connect(*csiEndpoint, metricsManager) ... // CSI服务就绪探测 ctrl.Probe(grpcClient, *operationTimeout) ... // 交互CSI实现获取名称 provisionerName, err := ctrl.GetDriverName(grpcClient, *operationTimeout) ... // in-tree-->out-tree if translator.IsMigratedCSIDriverByName(provisionerName) { // 迁移前的in-tree名称 supportsMigrationFromInTreePluginName, err = translator.GetInTreeNameFromCSIName(provisionerName) ... // 重新建立grpc连接 migratedGrpcClient, err := ctrl.Connect(*csiEndpoint, metricsManager) ... grpcClient.Close() grpcClient = migratedGrpcClient // CSI服务就绪探测 ctrl.Probe(grpcClient, *operationTimeout) ... } ... // 获取CSI能力 pluginCapabilities, controllerCapabilities, err := ctrl.GetDriverCapabilities(grpcClient, *operationTimeout) ... // Listers scLister := factory.Storage().V1().StorageClasses().Lister() claimLister := factory.Core().V1().PersistentVolumeClaims().Lister() ... // CSI服务支持Attacher/Detacher if controllerCapabilities[csi.ControllerServiceCapability_RPC_PUBLISH_UNPUBLISH_VOLUME] { vaLister = factory.Storage().V1().VolumeAttachments().Lister() } ... // ds部署模式,Pod仅管理所在node if *enableNodeDeployment { nodeDeployment = &ctrl.NodeDeployment{ NodeName: node, ClaimInformer: factory.Core().V1().PersistentVolumeClaims(), // PVC监听 ImmediateBinding: *nodeDeploymentImmediateBinding, // PVC绑定模式(立即/延迟) BaseDelay: *nodeDeploymentBaseDelay, // 20s基础延迟 MaxDelay: *nodeDeploymentMaxDelay, // 60s最大延迟 } // 获取CSI服务node信息(nodeID/最大卷数/可能的拓扑) nodeInfo, err := ctrl.GetNodeInfo(grpcClient, *operationTimeout) ... nodeDeployment.NodeInfo = *nodeInfo } ... // CSI Driver创建卷考虑Topology if ctrl.SupportsTopology(pluginCapabilities) { // nodeDeploy模式 if nodeDeployment != nil { // 伪造一个CSINode对象 csiNode := &storagev1.CSINode{ ObjectMeta: metav1.ObjectMeta{ Name: nodeDeployment.NodeName, }, Spec: storagev1.CSINodeSpec{ Drivers: []storagev1.CSINodeDriver{ { Name: provisionerName, NodeID: nodeDeployment.NodeInfo.NodeId, }, }, }, } // 伪造一个node node := &v1.Node{ ObjectMeta: metav1.ObjectMeta{ Name: nodeDeployment.NodeName, }, } // 向CSINode和Node注入拓扑信息 if nodeDeployment.NodeInfo.AccessibleTopology != nil { for key := range nodeDeployment.NodeInfo.AccessibleTopology.Segments { csiNode.Spec.Drivers[0].TopologyKeys = append(csiNode.Spec.Drivers[0].TopologyKeys, key) } node.Labels = nodeDeployment.NodeInfo.AccessibleTopology.Segments } // 虚拟informer,不会List/Watch底层资源 stoppedFactory := informers.NewSharedInformerFactory(clientset, 1000*time.Hour) csiNodes := stoppedFactory.Storage().V1().CSINodes() nodes := stoppedFactory.Core().V1().Nodes() // 伪造的CSINode和Node注入Store csiNodes.Informer().GetStore().Add(csiNode) nodes.Informer().GetStore().Add(node) csiNodeLister = csiNodes.Lister() nodeLister = nodes.Lister() // 非nodeDeploy模式 } else { // 监听及缓存底层CSINode和Node csiNodeLister = factory.Storage().V1().CSINodes().Lister() nodeLister = factory.Core().V1().Nodes().Lister() } } ... // 启用跨命名空间Volume引用 if utilfeature.DefaultFeatureGate.Enabled(features.CrossNamespaceVolumeDataSource) { // 监听及缓存ReferenceGrant对象(引用限制) gatewayFactory = gatewayInformers.NewSharedInformerFactory(gatewayClient, 1h) referenceGrants := gatewayFactory.Gateway().V1beta1().ReferenceGrants() referenceGrantLister = referenceGrants.Lister() } // PersistentVolumeClaims informer rateLimiter := workqueue.NewItemExponentialFailureRateLimiter(1s, 5s) claimQueue := workqueue.NewNamedRateLimitingQueue(rateLimiter, "claims") claimInformer := factory.Core().V1().PersistentVolumeClaims().Informer() ... // create the provisioner csiProvisioner := ctrl.NewCSIProvisioner(...) ... // 启用容量检测 if *enableCapacity { ... if *capacityOwnerrefLevel >= 0 { podName := os.Getenv("POD_NAME") ... // 基于level递归找owner controller, err = owner.Lookup(config, namespace, podName, schema.GroupVersionKind{ Group: "", Version: "v1", Kind: "Pod" }, *capacityOwnerrefLevel) ... } ... // 非nodeDeploy模式 if nodeDeployment == nil { // 监听CSINode/Node,获取运行CSI Driver的Node TopologyKeys topologyInformer = topology.NewNodeTopology(provisionerName, clientset, factory.Core().V1().Nodes(), factory.Storage().V1().CSINodes(), workqueue.NewNamedRateLimitingQueue(rateLimiter, "csitopology")) // nodeDeploy模式 } else { ... // 收集CSI Driver上报的节点拓扑信息 if nodeDeployment.NodeInfo.AccessibleTopology != nil { for key, value := range nodeDeployment.NodeInfo.AccessibleTopology.Segments { segment = append(segment, topology.SegmentEntry{Key: key, Value: value}) } } // 基于CSI上报的构建topologyInformer topologyInformer = topology.NewFixedNodeTopology(&segment) } // 启动topologyInformer go topologyInformer.RunWorker(ctx) // external-provisioner身份 managedByID := "external-provisioner" // nodeDeploy模式 if *enableNodeDeployment { // 拼接node标记保证唯一 managedByID = getNameWithMaxLength(managedByID, node, validation.DNS1035LabelMaxLength) } ... // 先监听V1 CSIStorageCapacity API资源 clientFactory := capacity.NewV1ClientFactory(clientset) cInformer := fn.Storage().V1().CSIStorageCapacities() // 构造一个无效的V1 CSIStorageCapacity invalidCapacity := &storagev1.CSIStorageCapacity{ ObjectMeta: metav1.ObjectMeta{ Name: "#%123-invalid-name", }, } // 尝试创建V1 invalidCapacity检查合法性 createdCapacity, err := clientset.StorageV1().CSIStorageCapacities(namespace).Create(ctx, invalidCapacity, metav1.CreateOptions{}) switch { case err == nil: // succeed case apierrors.IsNotFound(err): // 不支持,转为v1beta1 clientFactory = capacity.NewV1beta1ClientFactory(clientset) cInformer = capacity.NewV1beta1InformerBridge(fn.Storage().V1beta1().CSIStorageCapacities()) ... } // 实例化capacityController capacityController = capacity.NewCentralCapacityController(...) ... // 构造容量包装的csiProvisioner csiProvisioner = capacity.NewProvisionWrapper(csiProvisioner, capacityController) } // 实例化provision Controller provisionController = controller.NewProvisionController( clientset, provisionerName, csiProvisioner, provisionerOptions...) // 实例化pvc controller csiClaimController := ctrl.NewCloningProtectionController( clientset, claimLister, claimInformer, claimQueue, controllerCapabilities) ... run := func(ctx context.Context) { // 开启PVC/StorageClass/Node/CSINode/VolumeAttachment监听同步 factory.Start(ctx.Done()) if fn != nil { // 开始CSIStorageCapacity监听同步 fn.Start(ctx.Done()) } ... // 开启跨命名空间volume引用 if utilfeature.DefaultFeatureGate.Enabled(features.CrossNamespaceVolumeDataSource) { // 启动ReferenceGrant监听同步 gatewayFactory.Start(ctx.Done()) ... } // 开启CSIStorageCapacity监听同步 if capacityController != nil { go capacityController.Run(ctx, int(*capacityThreads)) } // 开启Volume克隆 if csiClaimController != nil { go csiClaimController.Run(ctx, int(*finalizerThreads)) } // 开启provisioner协调 provisionController.Run(ctx) } // 无需选举 if !*enableLeaderElection { run(ctx) // 选举 } else { // 锁名 lockName := strings.Replace(provisionerName, "/", "-", -1) ... // 实例化leaderElection le := leaderelection.NewLeaderElection(leClientset, lockName, run) ... // lease所在namespace if *leaderElectionNamespace != "" { le.WithNamespace(*leaderElectionNamespace) } // 租约15s le.WithLeaseDuration(*leaderElectionLeaseDuration) // 续期窗口10s le.WithRenewDeadline(*leaderElectionRenewDeadline) // 重试间隔5s le.WithRetryPeriod(*leaderElectionRetryPeriod) // 身份 le.WithIdentity(identity) // 选举运行 le.Run() ... } }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
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288注意
mainLoop会初始化CSI Driver连接,构造各类资源informer/lister监听及缓存,实例化干活的controller
# 2.topology
# 2.1.初始化
topologyInformer初始化分为两种形式,基于informer监听的拓扑维护和基于CSI Driver限制的拓扑维护,后者交互CSI Driver同步。// returns an informer that synthesizes storage topology segments based on accessible topology that each CSI // driver node instance reports. func NewNodeTopology(...) Informer { // 实例化nodeTopology nt := &nodeTopology{ driverName: driverName, client: client, nodeInformer: nodeInformer, csiNodeInformer: csiNodeInformer, queue: queue, } // node事件监听 nodeHandler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { node, ok := obj.(*v1.Node) ... queue.Add("") }, UpdateFunc: func(oldObj interface{}, newObj interface{}) { oldNode, ok := oldObj.(*v1.Node) ... newNode, ok := newObj.(*v1.Node) ... // label无变化 if reflect.DeepEqual(oldNode.Labels, newNode.Labels) { return } queue.Add("") }, DeleteFunc: func(obj interface{}) { ... node, ok := obj.(*v1.Node) ... queue.Add("") }, } nodeInformer.Informer().AddEventHandler(nodeHandler) // csiNode事件监听 csiNodeHandler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { csiNode, ok := obj.(*storagev1.CSINode) ... queue.Add("") }, UpdateFunc: func(oldObj interface{}, newObj interface{}) { oldCSINode, ok := oldObj.(*storagev1.CSINode) ... newCSINode, ok := newObj.(*storagev1.CSINode) ... // 相关的TopologyKeys无变化 if reflect.DeepEqual(oldKeys, newKeys) { return } queue.Add("") }, DeleteFunc: func(obj interface{}) { ... csiNode, ok := obj.(*storagev1.CSINode) ... queue.Add("") }, } csiNodeInformer.Informer().AddEventHandler(csiNodeHandler) return nt } // creates topology informer for a driver with a fixed topology segment. func NewFixedNodeTopology(segment *Segment) Informer { return NewMock(segment) }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
补充
topologyInformer资源监听相对简单,相当于给queue入队一个可合并的token,由worker将多次变更合并处理
# 2.2.runworker
topologyInformer.runWorker()负责消费队列项,基于CSINode/Node和已缓存的TopologyKeys更新内部状态,维护节点的拓扑信息。func (nt *nodeTopology) RunWorker(ctx context.Context) { ... for nt.processNextWorkItem(ctx) { } } func (nt *nodeTopology) processNextWorkItem(ctx context.Context) bool { obj, shutdown := nt.queue.Get() ... defer nt.queue.Done(obj) nt.sync(ctx) return true } func (nt *nodeTopology) sync(ctx context.Context) { // 已缓存的TopologyKey segments := nt.List() ... // 假设旧的均要删除 for _, segment := range segments { removalCandidates[segment] = true } // 获取CSINode csiNodes, err := nt.csiNodeInformer.Lister().List(labels.Everything()) ... node: for _, csiNode := range csiNodes { // driver匹配的CSINode TopologyKey topologyKeys := nt.driverTopologyKeys(csiNode) if topologyKeys == nil { continue } // 获取node对象 node, err := nt.nodeInformer.Lister().Get(csiNode.Name) ... // 排序 sort.Strings(topologyKeys) // node.Labels匹配TopologyKey for _, key := range topologyKeys { value, ok := node.Labels[key] // 有一个不匹配就跳过 if !ok { continue node } newSegment = append(newSegment, SegmentEntry{key, value}) } // 对账旧的 for _, segment := range segments { // segment还在用 if newSegment.Compare(*segment) == 0 { // 标记segment未删除 removalCandidates[segment] = false // 补充到已存在集合 existingSegments = append(existingSegments, segment) continue node } } // 对账新增的 for _, segment := range addedSegments { // 同一轮不同node生成相同的Topolog就跳过 if newSegment.Compare(*segment) == 0 { // We already discovered this new segment. continue node } } // 更新segment集合 addedSegments = append(addedSegments, &newSegment) existingSegments = append(existingSegments, &newSegment) } ... nt.segments = existingSegments ... // 更新removedSegments集合 for segment, wasRemoved := range removalCandidates { if wasRemoved { removedSegments = append(removedSegments, segment) } } // 存在变动 if len(addedSegments) > 0 || len(removedSegments) > 0 { // 外部注册回调 for _, cb := range callbacks { cb(addedSegments, removedSegments) } } }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
注意
runworker基于现有的segments和CSINode匹配的segments计算差异,出现新增或删除segment执行回调(onTopologyChanges)
# 2.3.callback
capacityController初始化会向topologyInformer注册callback,本质上注册的是c.onTopologyChanges()处理segments变化。// called by the topology informer as callback. func (c *Controller) onTopologyChanges(added []*topology.Segment, removed []*topology.Segment) { // 获取storageClass列表 storageclasses, err := c.scInformer.Lister().List(labels.Everything()) ... for _, sc := range storageclasses { // 仅关心CSIDriver匹配的 if sc.Provisioner != c.driverName { continue } // 绑定语义冲突(控制器不支持立即绑定,storageClass要求立即绑定) if !c.immediateBinding && *sc.VolumeBindingMode == storagev1.VolumeBindingImmediate { continue } // add项入队 for _, segment := range added { c.addWorkItem(segment, sc) } // remove项入队 for _, segment := range removed { c.removeWorkItem(segment, sc) } } } // ensures that there is an item in c.capacities. It must be called while holding c.capacitiesLock! func (c *Controller) addWorkItem(segment *topology.Segment, sc *storagev1.StorageClass) { item := workItem{ segment: segment, storageClassName: sc.Name, } // 先认领 _, found := c.capacities[item] if !found { c.capacities[item] = nil } // 推到capacityQueue c.queue.Add(item) } // ensures that the item gets removed from c.capacities. It must be called while holding c.capacitiesLock! func (c *Controller) removeWorkItem(segment *topology.Segment, sc *storagev1.StorageClass) { item := workItem{ segment: segment, storageClassName: sc.Name, } // 未认领过 capacity, found := c.capacities[item] if !found { return } // 解注册 delete(c.capacities, item) // 未注册过capacity对象 if capacity == nil { return } // 推到capacityQueue c.queue.Add(capacity) }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
注意
onTopologyChanges基于segment和匹配的storageClass注册或解注册capacity对象及入队影响容量计算
# 3.capacity
# 3.1.初始化
NewCentralCapacityController()会实例化capacityController,注册StorageClass回调及向topologyInformer注册callback。// creates a new controller for CSIStorageCapacity objects. func NewCentralCapacityController(...) *Controller { // 实例化capacityController c := &Controller{...} handler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { sc, ok := obj.(*storagev1.StorageClass) ... c.onSCAddOrUpdate(sc) }, UpdateFunc: func(_ interface{}, newObj interface{}) { sc, ok := newObj.(*storagev1.StorageClass) ... c.onSCAddOrUpdate(sc) }, DeleteFunc: func(obj interface{}) { ... sc, ok := obj.(*storagev1.StorageClass) ... c.onSCDelete(sc) }, } // 注册informer回调 c.scInformer.Informer().AddEventHandler(handler) // 向topologyInformer注册回调触发容量调整 c.topologyInformer.AddCallback(c.onTopologyChanges) // 仅监听当前NS的Label Driver匹配的CSIStorageCapacity cInformer.Informer() return c } // onSCAddOrUpdate is called for add or update events by the storage class listener. func (c *Controller) onSCAddOrUpdate(sc *storagev1.StorageClass) { // provisioner匹配检查 if sc.Provisioner != c.driverName { return } // 模式兼容检查 if !c.immediateBinding && sc.VolumeBindingMode != nil && *sc.VolumeBindingMode == VolumeBindingImmediate { return } // 获取topologyInformer缓存的segment segments := c.topologyInformer.List() ... // 注册capacity及capacity入队 for _, segment := range segments { c.addWorkItem(segment, sc) } } // onSCDelete is called for delete events by the storage class listener. func (c *Controller) onSCDelete(sc *storagev1.StorageClass) { if sc.Provisioner != c.driverName { return } // 获取topologyInformer缓存的segment segments := c.topologyInformer.List() ... // 解注册capacity及capacity入队 for _, segment := range segments { c.removeWorkItem(segment, sc) } }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
注意
capacityController会监听storageClass变化,配合topologyInformer入队capacity相关资源触发容量调整
# 3.2.prepare
capacityController.prepare()主动推送全量segment事件,补充CSIStorageCapacity事件监听以更新capacities缓存及队列任务项。// takes a read-only CSIStorageCapacity object and either remembers the pointer to it. func (c *Controller) onCAddOrUpdate(ctx context.Context, capacity *storagev1.CSIStorageCapacity) { // label driver及manageByID未匹配 if !c.isManaged(capacity) { ... for item, capacity2 := range c.capacities { // // 注册过 if capacity2 != nil && capacity2.UID == capacity.UID { // 置空 c.capacities[item] = nil // 入队重建 c.queue.Add(item) } } return } ... for item, capacity2 := range c.capacities { // 注册过 if capacity2 != nil && capacity2.UID == capacity.UID { // 更新缓存 c.capacities[item] = capacity return } // 认领过+SCName匹配及label匹配 if capacity2 == nil && item.equals(capacity) { // 注册 c.capacities[item] = capacity return } } // 入队重建 c.queue.Add(capacity) } func (c *Controller) onCDelete(ctx context.Context, capacity *storagev1.CSIStorageCapacity) { ... for item, capacity2 := range c.capacities { // 注册过 if capacity2 != nil && capacity2.UID == capacity.UID { // 置空 c.capacities[item] = nil // 入队重建 c.queue.Add(item) return } } } func (c *Controller) prepare(ctx context.Context) { // Wait for topology and storage class informer sync. if !cache.WaitForCacheSync(..., topologyInformer.HasSynced, scInformer.HasSynced, cInformer.HasSynced) { return } // 主动推送一次全量事件 c.onTopologyChanges(c.topologyInformer.List(), nil) ... // CSIStorageCapacity监听回调 handler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { csc, ok := obj.(*storagev1.CSIStorageCapacity) ... c.onCAddOrUpdate(ctx, csc) }, UpdateFunc: func(_ interface{}, newObj interface{}) { csc, ok := newObj.(*storagev1.CSIStorageCapacity) ... c.onCAddOrUpdate(ctx, csc) }, DeleteFunc: func(obj interface{}) { csc, ok := obj.(*storagev1.CSIStorageCapacity) ... c.onCDelete(ctx, csc) }, } c.cInformer.Informer().AddEventHandler(handler) // CSIStorageCapacity缓存 capacities, err := c.cInformer.Lister().List(labels.Everything()) ... // 推送capacity for _, capacity := range capacities { c.onCAddOrUpdate(ctx, capacity) } }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
注意
c.prepare()主要做一些准备工作,诸如首次启动全量事件推送及capacity对象监听回调注册
# 3.3.run入口
capacityController.Run()调用prepare()进行初始化及准备工作,启动消费worker和周期轮询worker异步处理容量调整事件。// Run is a main Controller handler func (c *Controller) Run(ctx context.Context, threadiness int) { ... // 启动前准备工作 c.prepare(ctx) // 消费worker并发数,默认为1 for i := 0; i < threadiness; i++ { go wait.UntilWithContext(ctx, func(ctx context.Context) { c.runWorker(ctx) }, time.Second) } // 间隔1min推送一次已有事件 go wait.UntilWithContext(ctx, func(ctx context.Context) { c.pollCapacities() }, c.pollPeriod) <-ctx.Done() } // pollCapacities must be called periodically to detect when the underlying storage capacity has changed. func (c *Controller) pollCapacities() { ... // workItem入队 for item := range c.capacities { c.queue.Add(item) } } func (c *Controller) runWorker(ctx context.Context) { for c.processNextWorkItem(ctx) { } } // processNextWorkItem processes items from queue. func (c *Controller) processNextWorkItem(ctx context.Context) bool { obj, shutdown := c.queue.Get() ... err := func() error { defer c.queue.Done(obj) switch obj := obj.(type) { case workItem: // 创建或更新capacity对象 return c.syncCapacity(ctx, obj) case *storagev1.CSIStorageCapacity: // 清理capacity对象 return c.deleteCapacity(ctx, obj) ... } return nil }() if err != nil { c.queue.AddRateLimited(obj) } else { c.queue.Forget(obj) } 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
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
注意
pollcapacity间隔1min将已有workItem入队触发一次全量检查,runWorker消费queue执行容量同步或删除
# 3.4.sync
c.syncCapacity()会获取关联的storageClass对象,交互driver获取存储容量,基于capacity注册状态创建或更新对象。// syncCapacity gets the capacity and then updates or creates the object. func (c *Controller) syncCapacity(ctx context.Context, item workItem) error { ... // 获取注册的capacity capacity, found := c.capacities[item] ... // 未注册过 if !found { return nil } // 获取对应storageClass对象 sc, err := c.scInformer.Lister().Get(item.storageClassName) ... // 获取CSI Driver存储容量 resp, err := c.csiController.GetCapacity(syncCtx, req) ... quantity := resource.NewQuantity(resp.AvailableCapacity, resource.BinarySI) ... if resp.MaximumVolumeSize != nil { maximumVolumeSize = resource.NewQuantity(resp.MaximumVolumeSize.Value, resource.BinarySI) } // 认领未注册过capacity if capacity == nil { // 构造capacity对象 capacity = &storagev1.CSIStorageCapacity{ ObjectMeta: metav1.ObjectMeta{ GenerateName: "csisc-", Labels: map[string]string{ DriverNameLabel: c.driverName, ManagedByLabel: c.managedByID, }, }, StorageClassName: item.storageClassName, NodeTopology: item.segment.GetLabelSelector(), Capacity: quantity, MaximumVolumeSize: maximumVolumeSize, } // 设置owner if c.owner != nil { capacity.OwnerReferences = []metav1.OwnerReference{*c.owner} } ... // 创建 capacity, err = c.clientFactory(c.ownerNamespace).Create(ctx, capacity, metav1.CreateOptions{}) ... // capacity容量未变且owner未变 } else if capacity.Capacity.Value() == quantity.Value() && sizesAreEqual(capacity.MaximumVolumeSize, maximumVolumeSize) && (c.owner == nil || c.isOwnedByUs(capacity)) { return nil } else { // 更新capacity capacity := capacity.DeepCopy() capacity.Capacity = quantity capacity.MaximumVolumeSize = maximumVolumeSize if c.owner != nil && !c.isOwnedByUs(capacity) { capacity.OwnerReferences = append(capacity.OwnerReferences, *c.owner) } ... capacity, err = c.clientFactory(capacity.Namespace).Update(ctx, capacity, metav1.UpdateOptions{}) ... } 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
注意
syncCapacity会基于capacity注册状态创建或更新capacity对象及更新capacities缓存
# 4.claim
# 4.1.run入口
csiClaimController.Run()负责注册PVC监听回调,启动异步worker消费队列项,驱动删除状态PVC对应存储的释放及回收处理。// enqueueClaimUpdate takes a PVC obj and stores it into the claim work queue. func (p *CloningProtectionController) enqueueClaimUpdate(ctx context.Context, obj interface{}) { new, ok := obj.(*v1.PersistentVolumeClaim) ... // 非删除状态 if new.DeletionTimestamp == nil { return } // Name/Namespace key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj) ... p.claimQueue.Add(key) } // Run is a main CloningProtectionController handler func (p *CloningProtectionController) Run(ctx context.Context, threadiness int) { ... // PVC监听回调 claimHandler := cache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { p.enqueueClaimUpdate(ctx, obj) }, UpdateFunc: func(_ interface{}, newObj interface{}) { p.enqueueClaimUpdate(ctx, newObj) }, } p.claimInformer.AddEventHandlerWithResyncPeriod(claimHandler, controller.DefaultResyncPeriod) // 启动PVC worker for i := 0; i < threadiness; i++ { go wait.Until(func() { p.runClaimWorker(ctx) }, 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注意
claim队列仅关心处于删除状态的PVC,驱动关联存储资源的释放和回收
# 4.2.claimworker
p.runClaimWorker()会消费队列的PVC,检查克隆进度,PVC克隆结束会清理克隆相关的finalizer以确保PVC正常删除。func (p *CloningProtectionController) runClaimWorker(ctx context.Context) { for p.processNextClaimWorkItem(ctx) { } } // processNextClaimWorkItem processes items from claimQueue func (p *CloningProtectionController) processNextClaimWorkItem(ctx context.Context) bool { obj, shutdown := p.claimQueue.Get() ... err := func(obj interface{}) error { defer p.claimQueue.Done(obj) ... // 执行同步 if err := p.syncClaimHandler(ctx, key); err != nil { p.claimQueue.AddRateLimited(obj) } else { p.claimQueue.Forget(obj) } return nil }(obj) ... return true } // syncClaimHandler gets the claim from informer's cache then calls syncClaim func (p *CloningProtectionController) syncClaimHandler(ctx context.Context, key string) error { namespace, name, err := cache.SplitMetaNamespaceKey(key) ... // 获取PVC对象 claim, err := p.claimLister.PersistentVolumeClaims(namespace).Get(name) ... // 同步PVC状态 return p.syncClaim(ctx, claim) } // syncClaim removes finalizers from a PVC, when cloning is finished func (p *CloningProtectionController) syncClaim(ctx context.Context, claim *v1.PersistentVolumeClaim) error { // 没有关心的finalizer if !checkFinalizer(claim, pvcCloneFinalizer) { return nil } // 获取PVC列表 pvcList, err := p.claimLister.PersistentVolumeClaims(claim.Namespace).List(labels.Everything()) ... // Check for pvc state with DataSource pointing to claim for _, pvc := range pvcList { // 不属于克隆PVC if pvc.Spec.DataSource == nil { continue } // 克隆PVC未完成 if pvc.DataSource.Kind == pvcKind && pvc.DataSource.Name == claim.Name && pvc.Phase == v1.ClaimPending { return fmt.Errorf("PVC '%s' is in 'Pending' state, cloning in progress", pvc.Name) } } ... // 移除关心的finalizer for _, finalizer := range claim.ObjectMeta.Finalizers { if finalizer != pvcCloneFinalizer { finalizers = append(finalizers, finalizer) } } claim.ObjectMeta.Finalizers = finalizers // 更新PVC p.client.CoreV1().PersistentVolumeClaims(claim.Namespace).Update(ctx, claim, metav1.UpdateOptions{}) ... 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
注意
claimworker会检查deletingPVC关联的clonePVC完成情况,满足条件会摘除clone finalizer以确保PVC正常删除