• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

kubeovn / kube-ovn / 18735454196

23 Oct 2025 02:16AM UTC coverage: 21.151% (-0.002%) from 21.153%
18735454196

push

github

web-flow
add lock for gateway exec operations (#5817)

Signed-off-by: Mengxin Liu <liumengxinfly@gmail.com>

0 of 18 new or added lines in 2 files covered. (0.0%)

10737 of 50764 relevant lines covered (21.15%)

0.25 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

1.2
/pkg/controller/controller.go
1
package controller
2

3
import (
4
        "context"
5
        "fmt"
6
        "runtime"
7
        "strings"
8
        "time"
9

10
        netAttach "github.com/k8snetworkplumbingwg/network-attachment-definition-client/pkg/client/informers/externalversions"
11
        netAttachv1 "github.com/k8snetworkplumbingwg/network-attachment-definition-client/pkg/client/listers/k8s.cni.cncf.io/v1"
12
        "github.com/puzpuzpuz/xsync/v4"
13
        "golang.org/x/time/rate"
14
        corev1 "k8s.io/api/core/v1"
15
        metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
16
        "k8s.io/apimachinery/pkg/labels"
17
        utilruntime "k8s.io/apimachinery/pkg/util/runtime"
18
        "k8s.io/apimachinery/pkg/util/wait"
19
        kubeinformers "k8s.io/client-go/informers"
20
        "k8s.io/client-go/kubernetes/scheme"
21
        typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
22
        appsv1 "k8s.io/client-go/listers/apps/v1"
23
        certListerv1 "k8s.io/client-go/listers/certificates/v1"
24
        v1 "k8s.io/client-go/listers/core/v1"
25
        discoveryv1 "k8s.io/client-go/listers/discovery/v1"
26
        netv1 "k8s.io/client-go/listers/networking/v1"
27
        "k8s.io/client-go/tools/cache"
28
        "k8s.io/client-go/tools/record"
29
        "k8s.io/client-go/util/workqueue"
30
        "k8s.io/klog/v2"
31
        "k8s.io/utils/keymutex"
32
        v1alpha1 "sigs.k8s.io/network-policy-api/apis/v1alpha1"
33
        anpinformer "sigs.k8s.io/network-policy-api/pkg/client/informers/externalversions"
34
        anplister "sigs.k8s.io/network-policy-api/pkg/client/listers/apis/v1alpha1"
35

36
        "github.com/kubeovn/kube-ovn/pkg/informer"
37

38
        kubeovnv1 "github.com/kubeovn/kube-ovn/pkg/apis/kubeovn/v1"
39
        kubeovninformer "github.com/kubeovn/kube-ovn/pkg/client/informers/externalversions"
40
        kubeovnlister "github.com/kubeovn/kube-ovn/pkg/client/listers/kubeovn/v1"
41
        ovnipam "github.com/kubeovn/kube-ovn/pkg/ipam"
42
        "github.com/kubeovn/kube-ovn/pkg/ovs"
43
        "github.com/kubeovn/kube-ovn/pkg/util"
44
)
45

46
const controllerAgentName = "kube-ovn-controller"
47

48
const (
49
        logicalSwitchKey              = "ls"
50
        logicalRouterKey              = "lr"
51
        portGroupKey                  = "pg"
52
        networkPolicyKey              = "np"
53
        sgKey                         = "sg"
54
        associatedSgKeyPrefix         = "associated_sg_"
55
        sgsKey                        = "security_groups"
56
        u2oKey                        = "u2o"
57
        adminNetworkPolicyKey         = "anp"
58
        baselineAdminNetworkPolicyKey = "banp"
59
)
60

61
// Controller is kube-ovn main controller that watch ns/pod/node/svc/ep and operate ovn
62
type Controller struct {
63
        config *Configuration
64

65
        ipam           *ovnipam.IPAM
66
        namedPort      *NamedPort
67
        anpPrioNameMap map[int32]string
68
        anpNamePrioMap map[string]int32
69

70
        OVNNbClient ovs.NbClient
71
        OVNSbClient ovs.SbClient
72

73
        // ExternalGatewayType define external gateway type, centralized
74
        ExternalGatewayType string
75

76
        podsLister             v1.PodLister
77
        podsSynced             cache.InformerSynced
78
        addOrUpdatePodQueue    workqueue.TypedRateLimitingInterface[string]
79
        deletePodQueue         workqueue.TypedRateLimitingInterface[string]
80
        deletingPodObjMap      *xsync.Map[string, *corev1.Pod]
81
        deletingNodeObjMap     *xsync.Map[string, *corev1.Node]
82
        updatePodSecurityQueue workqueue.TypedRateLimitingInterface[string]
83
        podKeyMutex            keymutex.KeyMutex
84

85
        vpcsLister           kubeovnlister.VpcLister
86
        vpcSynced            cache.InformerSynced
87
        addOrUpdateVpcQueue  workqueue.TypedRateLimitingInterface[string]
88
        vpcLastPoliciesMap   *xsync.Map[string, string]
89
        delVpcQueue          workqueue.TypedRateLimitingInterface[*kubeovnv1.Vpc]
90
        updateVpcStatusQueue workqueue.TypedRateLimitingInterface[string]
91
        vpcKeyMutex          keymutex.KeyMutex
92

93
        vpcNatGatewayLister           kubeovnlister.VpcNatGatewayLister
94
        vpcNatGatewaySynced           cache.InformerSynced
95
        addOrUpdateVpcNatGatewayQueue workqueue.TypedRateLimitingInterface[string]
96
        delVpcNatGatewayQueue         workqueue.TypedRateLimitingInterface[string]
97
        initVpcNatGatewayQueue        workqueue.TypedRateLimitingInterface[string]
98
        updateVpcEipQueue             workqueue.TypedRateLimitingInterface[string]
99
        updateVpcFloatingIPQueue      workqueue.TypedRateLimitingInterface[string]
100
        updateVpcDnatQueue            workqueue.TypedRateLimitingInterface[string]
101
        updateVpcSnatQueue            workqueue.TypedRateLimitingInterface[string]
102
        updateVpcSubnetQueue          workqueue.TypedRateLimitingInterface[string]
103
        vpcNatGwKeyMutex              keymutex.KeyMutex
104
        vpcNatGwExecKeyMutex          keymutex.KeyMutex
105

106
        vpcEgressGatewayLister           kubeovnlister.VpcEgressGatewayLister
107
        vpcEgressGatewaySynced           cache.InformerSynced
108
        addOrUpdateVpcEgressGatewayQueue workqueue.TypedRateLimitingInterface[string]
109
        delVpcEgressGatewayQueue         workqueue.TypedRateLimitingInterface[string]
110
        vpcEgressGatewayKeyMutex         keymutex.KeyMutex
111

112
        switchLBRuleLister      kubeovnlister.SwitchLBRuleLister
113
        switchLBRuleSynced      cache.InformerSynced
114
        addSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[string]
115
        updateSwitchLBRuleQueue workqueue.TypedRateLimitingInterface[*SlrInfo]
116
        delSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[*SlrInfo]
117

118
        vpcDNSLister           kubeovnlister.VpcDnsLister
119
        vpcDNSSynced           cache.InformerSynced
120
        addOrUpdateVpcDNSQueue workqueue.TypedRateLimitingInterface[string]
121
        delVpcDNSQueue         workqueue.TypedRateLimitingInterface[string]
122

123
        subnetsLister           kubeovnlister.SubnetLister
124
        subnetSynced            cache.InformerSynced
125
        addOrUpdateSubnetQueue  workqueue.TypedRateLimitingInterface[string]
126
        subnetLastVpcNameMap    *xsync.Map[string, string]
127
        deleteSubnetQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.Subnet]
128
        updateSubnetStatusQueue workqueue.TypedRateLimitingInterface[string]
129
        syncVirtualPortsQueue   workqueue.TypedRateLimitingInterface[string]
130
        subnetKeyMutex          keymutex.KeyMutex
131

132
        ippoolLister            kubeovnlister.IPPoolLister
133
        ippoolSynced            cache.InformerSynced
134
        addOrUpdateIPPoolQueue  workqueue.TypedRateLimitingInterface[string]
135
        updateIPPoolStatusQueue workqueue.TypedRateLimitingInterface[string]
136
        deleteIPPoolQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.IPPool]
137
        ippoolKeyMutex          keymutex.KeyMutex
138

139
        ipsLister     kubeovnlister.IPLister
140
        ipSynced      cache.InformerSynced
141
        addIPQueue    workqueue.TypedRateLimitingInterface[string]
142
        updateIPQueue workqueue.TypedRateLimitingInterface[string]
143
        delIPQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.IP]
144

145
        virtualIpsLister          kubeovnlister.VipLister
146
        virtualIpsSynced          cache.InformerSynced
147
        addVirtualIPQueue         workqueue.TypedRateLimitingInterface[string]
148
        updateVirtualIPQueue      workqueue.TypedRateLimitingInterface[string]
149
        updateVirtualParentsQueue workqueue.TypedRateLimitingInterface[string]
150
        delVirtualIPQueue         workqueue.TypedRateLimitingInterface[*kubeovnv1.Vip]
151

152
        iptablesEipsLister     kubeovnlister.IptablesEIPLister
153
        iptablesEipSynced      cache.InformerSynced
154
        addIptablesEipQueue    workqueue.TypedRateLimitingInterface[string]
155
        updateIptablesEipQueue workqueue.TypedRateLimitingInterface[string]
156
        resetIptablesEipQueue  workqueue.TypedRateLimitingInterface[string]
157
        delIptablesEipQueue    workqueue.TypedRateLimitingInterface[string]
158

159
        iptablesFipsLister     kubeovnlister.IptablesFIPRuleLister
160
        iptablesFipSynced      cache.InformerSynced
161
        addIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
162
        updateIptablesFipQueue workqueue.TypedRateLimitingInterface[string]
163
        delIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
164

165
        iptablesDnatRulesLister     kubeovnlister.IptablesDnatRuleLister
166
        iptablesDnatRuleSynced      cache.InformerSynced
167
        addIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
168
        updateIptablesDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
169
        delIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
170

171
        iptablesSnatRulesLister     kubeovnlister.IptablesSnatRuleLister
172
        iptablesSnatRuleSynced      cache.InformerSynced
173
        addIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
174
        updateIptablesSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
175
        delIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
176

177
        ovnEipsLister     kubeovnlister.OvnEipLister
178
        ovnEipSynced      cache.InformerSynced
179
        addOvnEipQueue    workqueue.TypedRateLimitingInterface[string]
180
        updateOvnEipQueue workqueue.TypedRateLimitingInterface[string]
181
        resetOvnEipQueue  workqueue.TypedRateLimitingInterface[string]
182
        delOvnEipQueue    workqueue.TypedRateLimitingInterface[string]
183

184
        ovnFipsLister     kubeovnlister.OvnFipLister
185
        ovnFipSynced      cache.InformerSynced
186
        addOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
187
        updateOvnFipQueue workqueue.TypedRateLimitingInterface[string]
188
        delOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
189

190
        ovnSnatRulesLister     kubeovnlister.OvnSnatRuleLister
191
        ovnSnatRuleSynced      cache.InformerSynced
192
        addOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
193
        updateOvnSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
194
        delOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
195

196
        ovnDnatRulesLister     kubeovnlister.OvnDnatRuleLister
197
        ovnDnatRuleSynced      cache.InformerSynced
198
        addOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
199
        updateOvnDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
200
        delOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
201

202
        providerNetworksLister kubeovnlister.ProviderNetworkLister
203
        providerNetworkSynced  cache.InformerSynced
204

205
        vlansLister     kubeovnlister.VlanLister
206
        vlanSynced      cache.InformerSynced
207
        addVlanQueue    workqueue.TypedRateLimitingInterface[string]
208
        delVlanQueue    workqueue.TypedRateLimitingInterface[string]
209
        updateVlanQueue workqueue.TypedRateLimitingInterface[string]
210
        vlanKeyMutex    keymutex.KeyMutex
211

212
        namespacesLister  v1.NamespaceLister
213
        namespacesSynced  cache.InformerSynced
214
        addNamespaceQueue workqueue.TypedRateLimitingInterface[string]
215
        nsKeyMutex        keymutex.KeyMutex
216

217
        nodesLister     v1.NodeLister
218
        nodesSynced     cache.InformerSynced
219
        addNodeQueue    workqueue.TypedRateLimitingInterface[string]
220
        updateNodeQueue workqueue.TypedRateLimitingInterface[string]
221
        deleteNodeQueue workqueue.TypedRateLimitingInterface[string]
222
        nodeKeyMutex    keymutex.KeyMutex
223

224
        servicesLister     v1.ServiceLister
225
        serviceSynced      cache.InformerSynced
226
        addServiceQueue    workqueue.TypedRateLimitingInterface[string]
227
        deleteServiceQueue workqueue.TypedRateLimitingInterface[*vpcService]
228
        updateServiceQueue workqueue.TypedRateLimitingInterface[*updateSvcObject]
229
        svcKeyMutex        keymutex.KeyMutex
230

231
        endpointSlicesLister          discoveryv1.EndpointSliceLister
232
        endpointSlicesSynced          cache.InformerSynced
233
        addOrUpdateEndpointSliceQueue workqueue.TypedRateLimitingInterface[string]
234
        epKeyMutex                    keymutex.KeyMutex
235

236
        deploymentsLister appsv1.DeploymentLister
237
        deploymentsSynced cache.InformerSynced
238

239
        npsLister     netv1.NetworkPolicyLister
240
        npsSynced     cache.InformerSynced
241
        updateNpQueue workqueue.TypedRateLimitingInterface[string]
242
        deleteNpQueue workqueue.TypedRateLimitingInterface[string]
243
        npKeyMutex    keymutex.KeyMutex
244

245
        sgsLister          kubeovnlister.SecurityGroupLister
246
        sgSynced           cache.InformerSynced
247
        addOrUpdateSgQueue workqueue.TypedRateLimitingInterface[string]
248
        delSgQueue         workqueue.TypedRateLimitingInterface[string]
249
        syncSgPortsQueue   workqueue.TypedRateLimitingInterface[string]
250
        sgKeyMutex         keymutex.KeyMutex
251

252
        qosPoliciesLister    kubeovnlister.QoSPolicyLister
253
        qosPolicySynced      cache.InformerSynced
254
        addQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
255
        updateQoSPolicyQueue workqueue.TypedRateLimitingInterface[string]
256
        delQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
257

258
        configMapsLister v1.ConfigMapLister
259
        configMapsSynced cache.InformerSynced
260

261
        anpsLister     anplister.AdminNetworkPolicyLister
262
        anpsSynced     cache.InformerSynced
263
        addAnpQueue    workqueue.TypedRateLimitingInterface[string]
264
        updateAnpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
265
        deleteAnpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.AdminNetworkPolicy]
266
        anpKeyMutex    keymutex.KeyMutex
267

268
        dnsNameResolversLister          kubeovnlister.DNSNameResolverLister
269
        dnsNameResolversSynced          cache.InformerSynced
270
        addOrUpdateDNSNameResolverQueue workqueue.TypedRateLimitingInterface[string]
271
        deleteDNSNameResolverQueue      workqueue.TypedRateLimitingInterface[*kubeovnv1.DNSNameResolver]
272

273
        banpsLister     anplister.BaselineAdminNetworkPolicyLister
274
        banpsSynced     cache.InformerSynced
275
        addBanpQueue    workqueue.TypedRateLimitingInterface[string]
276
        updateBanpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
277
        deleteBanpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.BaselineAdminNetworkPolicy]
278
        banpKeyMutex    keymutex.KeyMutex
279

280
        csrLister           certListerv1.CertificateSigningRequestLister
281
        csrSynced           cache.InformerSynced
282
        addOrUpdateCsrQueue workqueue.TypedRateLimitingInterface[string]
283

284
        addOrUpdateVMIMigrationQueue workqueue.TypedRateLimitingInterface[string]
285
        deleteVMQueue                workqueue.TypedRateLimitingInterface[string]
286
        kubevirtInformerFactory      informer.KubeVirtInformerFactory
287

288
        netAttachLister          netAttachv1.NetworkAttachmentDefinitionLister
289
        netAttachSynced          cache.InformerSynced
290
        netAttachInformerFactory netAttach.SharedInformerFactory
291

292
        recorder               record.EventRecorder
293
        informerFactory        kubeinformers.SharedInformerFactory
294
        cmInformerFactory      kubeinformers.SharedInformerFactory
295
        deployInformerFactory  kubeinformers.SharedInformerFactory
296
        kubeovnInformerFactory kubeovninformer.SharedInformerFactory
297
        anpInformerFactory     anpinformer.SharedInformerFactory
298

299
        // Database health check
300
        dbFailureCount int
301
}
302

303
func newTypedRateLimitingQueue[T comparable](name string, rateLimiter workqueue.TypedRateLimiter[T]) workqueue.TypedRateLimitingInterface[T] {
1✔
304
        if rateLimiter == nil {
2✔
305
                rateLimiter = workqueue.DefaultTypedControllerRateLimiter[T]()
1✔
306
        }
1✔
307
        return workqueue.NewTypedRateLimitingQueueWithConfig(rateLimiter, workqueue.TypedRateLimitingQueueConfig[T]{Name: name})
1✔
308
}
309

310
// Run creates and runs a new ovn controller
311
func Run(ctx context.Context, config *Configuration) {
×
312
        klog.V(4).Info("Creating event broadcaster")
×
313
        eventBroadcaster := record.NewBroadcasterWithCorrelatorOptions(record.CorrelatorOptions{BurstSize: 100})
×
314
        eventBroadcaster.StartLogging(klog.Infof)
×
315
        eventBroadcaster.StartRecordingToSink(&typedcorev1.EventSinkImpl{Interface: config.KubeFactoryClient.CoreV1().Events("")})
×
316
        recorder := eventBroadcaster.NewRecorder(scheme.Scheme, corev1.EventSource{Component: controllerAgentName})
×
317
        custCrdRateLimiter := workqueue.NewTypedMaxOfRateLimiter(
×
318
                workqueue.NewTypedItemExponentialFailureRateLimiter[string](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
319
                &workqueue.TypedBucketRateLimiter[string]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
320
        )
×
321

×
322
        selector, err := labels.Parse(util.VpcEgressGatewayLabel)
×
323
        if err != nil {
×
324
                util.LogFatalAndExit(err, "failed to create label selector for vpc egress gateway workload")
×
325
        }
×
326

327
        informerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
328
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
329
                        listOption.AllowWatchBookmarks = true
×
330
                }))
×
331
        cmInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
332
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
333
                        listOption.AllowWatchBookmarks = true
×
334
                }), kubeinformers.WithNamespace(config.PodNamespace))
×
335
        // deployment informer used to list/watch vpc egress gateway workloads
336
        deployInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
337
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
338
                        listOption.AllowWatchBookmarks = true
×
339
                        listOption.LabelSelector = selector.String()
×
340
                }))
×
341
        kubeovnInformerFactory := kubeovninformer.NewSharedInformerFactoryWithOptions(config.KubeOvnFactoryClient, 0,
×
342
                kubeovninformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
343
                        listOption.AllowWatchBookmarks = true
×
344
                }))
×
345
        anpInformerFactory := anpinformer.NewSharedInformerFactoryWithOptions(config.AnpClient, 0,
×
346
                anpinformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
347
                        listOption.AllowWatchBookmarks = true
×
348
                }))
×
349

350
        attachNetInformerFactory := netAttach.NewSharedInformerFactoryWithOptions(config.AttachNetClient, 0,
×
351
                netAttach.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
352
                        listOption.AllowWatchBookmarks = true
×
353
                }),
×
354
        )
355

356
        kubevirtInformerFactory := informer.NewKubeVirtInformerFactory(config.KubevirtClient.RestClient(), config.KubevirtClient, nil, util.KubevirtNamespace)
×
357

×
358
        vpcInformer := kubeovnInformerFactory.Kubeovn().V1().Vpcs()
×
359
        vpcNatGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcNatGateways()
×
360
        vpcEgressGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcEgressGateways()
×
361
        subnetInformer := kubeovnInformerFactory.Kubeovn().V1().Subnets()
×
362
        ippoolInformer := kubeovnInformerFactory.Kubeovn().V1().IPPools()
×
363
        ipInformer := kubeovnInformerFactory.Kubeovn().V1().IPs()
×
364
        virtualIPInformer := kubeovnInformerFactory.Kubeovn().V1().Vips()
×
365
        iptablesEipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesEIPs()
×
366
        iptablesFipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesFIPRules()
×
367
        iptablesDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesDnatRules()
×
368
        iptablesSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesSnatRules()
×
369
        vlanInformer := kubeovnInformerFactory.Kubeovn().V1().Vlans()
×
370
        providerNetworkInformer := kubeovnInformerFactory.Kubeovn().V1().ProviderNetworks()
×
371
        sgInformer := kubeovnInformerFactory.Kubeovn().V1().SecurityGroups()
×
372
        podInformer := informerFactory.Core().V1().Pods()
×
373
        namespaceInformer := informerFactory.Core().V1().Namespaces()
×
374
        nodeInformer := informerFactory.Core().V1().Nodes()
×
375
        serviceInformer := informerFactory.Core().V1().Services()
×
376
        endpointSliceInformer := informerFactory.Discovery().V1().EndpointSlices()
×
377
        deploymentInformer := deployInformerFactory.Apps().V1().Deployments()
×
378
        qosPolicyInformer := kubeovnInformerFactory.Kubeovn().V1().QoSPolicies()
×
379
        configMapInformer := cmInformerFactory.Core().V1().ConfigMaps()
×
380
        npInformer := informerFactory.Networking().V1().NetworkPolicies()
×
381
        switchLBRuleInformer := kubeovnInformerFactory.Kubeovn().V1().SwitchLBRules()
×
382
        vpcDNSInformer := kubeovnInformerFactory.Kubeovn().V1().VpcDnses()
×
383
        ovnEipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnEips()
×
384
        ovnFipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnFips()
×
385
        ovnSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnSnatRules()
×
386
        ovnDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnDnatRules()
×
387
        anpInformer := anpInformerFactory.Policy().V1alpha1().AdminNetworkPolicies()
×
388
        banpInformer := anpInformerFactory.Policy().V1alpha1().BaselineAdminNetworkPolicies()
×
389
        dnsNameResolverInformer := kubeovnInformerFactory.Kubeovn().V1().DNSNameResolvers()
×
390
        csrInformer := informerFactory.Certificates().V1().CertificateSigningRequests()
×
391
        netAttachInformer := attachNetInformerFactory.K8sCniCncfIo().V1().NetworkAttachmentDefinitions()
×
392

×
393
        numKeyLocks := max(runtime.NumCPU()*2, config.WorkerNum*2)
×
394
        controller := &Controller{
×
395
                config:             config,
×
396
                deletingPodObjMap:  xsync.NewMap[string, *corev1.Pod](),
×
397
                deletingNodeObjMap: xsync.NewMap[string, *corev1.Node](),
×
398
                ipam:               ovnipam.NewIPAM(),
×
399
                namedPort:          NewNamedPort(),
×
400

×
401
                vpcsLister:           vpcInformer.Lister(),
×
402
                vpcSynced:            vpcInformer.Informer().HasSynced,
×
403
                addOrUpdateVpcQueue:  newTypedRateLimitingQueue[string]("AddOrUpdateVpc", nil),
×
404
                vpcLastPoliciesMap:   xsync.NewMap[string, string](),
×
405
                delVpcQueue:          newTypedRateLimitingQueue[*kubeovnv1.Vpc]("DeleteVpc", nil),
×
406
                updateVpcStatusQueue: newTypedRateLimitingQueue[string]("UpdateVpcStatus", nil),
×
407
                vpcKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
408

×
NEW
409
                vpcNatGatewayLister:              vpcNatGatewayInformer.Lister(),
×
NEW
410
                vpcNatGatewaySynced:              vpcNatGatewayInformer.Informer().HasSynced,
×
NEW
411
                addOrUpdateVpcNatGatewayQueue:    newTypedRateLimitingQueue("AddOrUpdateVpcNatGw", custCrdRateLimiter),
×
NEW
412
                initVpcNatGatewayQueue:           newTypedRateLimitingQueue("InitVpcNatGw", custCrdRateLimiter),
×
NEW
413
                delVpcNatGatewayQueue:            newTypedRateLimitingQueue("DeleteVpcNatGw", custCrdRateLimiter),
×
NEW
414
                updateVpcEipQueue:                newTypedRateLimitingQueue("UpdateVpcEip", custCrdRateLimiter),
×
NEW
415
                updateVpcFloatingIPQueue:         newTypedRateLimitingQueue("UpdateVpcFloatingIp", custCrdRateLimiter),
×
NEW
416
                updateVpcDnatQueue:               newTypedRateLimitingQueue("UpdateVpcDnat", custCrdRateLimiter),
×
NEW
417
                updateVpcSnatQueue:               newTypedRateLimitingQueue("UpdateVpcSnat", custCrdRateLimiter),
×
NEW
418
                updateVpcSubnetQueue:             newTypedRateLimitingQueue("UpdateVpcSubnet", custCrdRateLimiter),
×
NEW
419
                vpcNatGwKeyMutex:                 keymutex.NewHashed(numKeyLocks),
×
NEW
420
                vpcNatGwExecKeyMutex:             keymutex.NewHashed(numKeyLocks),
×
421
                vpcEgressGatewayLister:           vpcEgressGatewayInformer.Lister(),
×
422
                vpcEgressGatewaySynced:           vpcEgressGatewayInformer.Informer().HasSynced,
×
423
                addOrUpdateVpcEgressGatewayQueue: newTypedRateLimitingQueue("AddOrUpdateVpcEgressGateway", custCrdRateLimiter),
×
424
                delVpcEgressGatewayQueue:         newTypedRateLimitingQueue("DeleteVpcEgressGateway", custCrdRateLimiter),
×
425
                vpcEgressGatewayKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
426

×
427
                subnetsLister:           subnetInformer.Lister(),
×
428
                subnetSynced:            subnetInformer.Informer().HasSynced,
×
429
                addOrUpdateSubnetQueue:  newTypedRateLimitingQueue[string]("AddSubnet", nil),
×
430
                subnetLastVpcNameMap:    xsync.NewMap[string, string](),
×
431
                deleteSubnetQueue:       newTypedRateLimitingQueue[*kubeovnv1.Subnet]("DeleteSubnet", nil),
×
432
                updateSubnetStatusQueue: newTypedRateLimitingQueue[string]("UpdateSubnetStatus", nil),
×
433
                syncVirtualPortsQueue:   newTypedRateLimitingQueue[string]("SyncVirtualPort", nil),
×
434
                subnetKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
435

×
436
                ippoolLister:            ippoolInformer.Lister(),
×
437
                ippoolSynced:            ippoolInformer.Informer().HasSynced,
×
438
                addOrUpdateIPPoolQueue:  newTypedRateLimitingQueue[string]("AddIPPool", nil),
×
439
                updateIPPoolStatusQueue: newTypedRateLimitingQueue[string]("UpdateIPPoolStatus", nil),
×
440
                deleteIPPoolQueue:       newTypedRateLimitingQueue[*kubeovnv1.IPPool]("DeleteIPPool", nil),
×
441
                ippoolKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
442

×
443
                ipsLister:     ipInformer.Lister(),
×
444
                ipSynced:      ipInformer.Informer().HasSynced,
×
445
                addIPQueue:    newTypedRateLimitingQueue[string]("AddIP", nil),
×
446
                updateIPQueue: newTypedRateLimitingQueue[string]("UpdateIP", nil),
×
447
                delIPQueue:    newTypedRateLimitingQueue[*kubeovnv1.IP]("DeleteIP", nil),
×
448

×
449
                virtualIpsLister:          virtualIPInformer.Lister(),
×
450
                virtualIpsSynced:          virtualIPInformer.Informer().HasSynced,
×
451
                addVirtualIPQueue:         newTypedRateLimitingQueue[string]("AddVirtualIP", nil),
×
452
                updateVirtualIPQueue:      newTypedRateLimitingQueue[string]("UpdateVirtualIP", nil),
×
453
                updateVirtualParentsQueue: newTypedRateLimitingQueue[string]("UpdateVirtualParents", nil),
×
454
                delVirtualIPQueue:         newTypedRateLimitingQueue[*kubeovnv1.Vip]("DeleteVirtualIP", nil),
×
455

×
456
                iptablesEipsLister:     iptablesEipInformer.Lister(),
×
457
                iptablesEipSynced:      iptablesEipInformer.Informer().HasSynced,
×
458
                addIptablesEipQueue:    newTypedRateLimitingQueue("AddIptablesEip", custCrdRateLimiter),
×
459
                updateIptablesEipQueue: newTypedRateLimitingQueue("UpdateIptablesEip", custCrdRateLimiter),
×
460
                resetIptablesEipQueue:  newTypedRateLimitingQueue("ResetIptablesEip", custCrdRateLimiter),
×
461
                delIptablesEipQueue:    newTypedRateLimitingQueue("DeleteIptablesEip", custCrdRateLimiter),
×
462

×
463
                iptablesFipsLister:     iptablesFipInformer.Lister(),
×
464
                iptablesFipSynced:      iptablesFipInformer.Informer().HasSynced,
×
465
                addIptablesFipQueue:    newTypedRateLimitingQueue("AddIptablesFip", custCrdRateLimiter),
×
466
                updateIptablesFipQueue: newTypedRateLimitingQueue("UpdateIptablesFip", custCrdRateLimiter),
×
467
                delIptablesFipQueue:    newTypedRateLimitingQueue("DeleteIptablesFip", custCrdRateLimiter),
×
468

×
469
                iptablesDnatRulesLister:     iptablesDnatRuleInformer.Lister(),
×
470
                iptablesDnatRuleSynced:      iptablesDnatRuleInformer.Informer().HasSynced,
×
471
                addIptablesDnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesDnatRule", custCrdRateLimiter),
×
472
                updateIptablesDnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesDnatRule", custCrdRateLimiter),
×
473
                delIptablesDnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesDnatRule", custCrdRateLimiter),
×
474

×
475
                iptablesSnatRulesLister:     iptablesSnatRuleInformer.Lister(),
×
476
                iptablesSnatRuleSynced:      iptablesSnatRuleInformer.Informer().HasSynced,
×
477
                addIptablesSnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesSnatRule", custCrdRateLimiter),
×
478
                updateIptablesSnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesSnatRule", custCrdRateLimiter),
×
479
                delIptablesSnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesSnatRule", custCrdRateLimiter),
×
480

×
481
                vlansLister:     vlanInformer.Lister(),
×
482
                vlanSynced:      vlanInformer.Informer().HasSynced,
×
483
                addVlanQueue:    newTypedRateLimitingQueue[string]("AddVlan", nil),
×
484
                delVlanQueue:    newTypedRateLimitingQueue[string]("DeleteVlan", nil),
×
485
                updateVlanQueue: newTypedRateLimitingQueue[string]("UpdateVlan", nil),
×
486
                vlanKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
487

×
488
                providerNetworksLister: providerNetworkInformer.Lister(),
×
489
                providerNetworkSynced:  providerNetworkInformer.Informer().HasSynced,
×
490

×
491
                podsLister:          podInformer.Lister(),
×
492
                podsSynced:          podInformer.Informer().HasSynced,
×
493
                addOrUpdatePodQueue: newTypedRateLimitingQueue[string]("AddOrUpdatePod", nil),
×
494
                deletePodQueue: workqueue.NewTypedRateLimitingQueueWithConfig(
×
495
                        workqueue.DefaultTypedControllerRateLimiter[string](),
×
496
                        workqueue.TypedRateLimitingQueueConfig[string]{
×
497
                                Name:          "DeletePod",
×
498
                                DelayingQueue: workqueue.NewTypedDelayingQueue[string](),
×
499
                        },
×
500
                ),
×
501
                updatePodSecurityQueue: newTypedRateLimitingQueue[string]("UpdatePodSecurity", nil),
×
502
                podKeyMutex:            keymutex.NewHashed(numKeyLocks),
×
503

×
504
                namespacesLister:  namespaceInformer.Lister(),
×
505
                namespacesSynced:  namespaceInformer.Informer().HasSynced,
×
506
                addNamespaceQueue: newTypedRateLimitingQueue[string]("AddNamespace", nil),
×
507
                nsKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
508

×
509
                nodesLister:     nodeInformer.Lister(),
×
510
                nodesSynced:     nodeInformer.Informer().HasSynced,
×
511
                addNodeQueue:    newTypedRateLimitingQueue[string]("AddNode", nil),
×
512
                updateNodeQueue: newTypedRateLimitingQueue[string]("UpdateNode", nil),
×
513
                deleteNodeQueue: newTypedRateLimitingQueue[string]("DeleteNode", nil),
×
514
                nodeKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
515

×
516
                servicesLister:     serviceInformer.Lister(),
×
517
                serviceSynced:      serviceInformer.Informer().HasSynced,
×
518
                addServiceQueue:    newTypedRateLimitingQueue[string]("AddService", nil),
×
519
                deleteServiceQueue: newTypedRateLimitingQueue[*vpcService]("DeleteService", nil),
×
520
                updateServiceQueue: newTypedRateLimitingQueue[*updateSvcObject]("UpdateService", nil),
×
521
                svcKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
522

×
523
                endpointSlicesLister:          endpointSliceInformer.Lister(),
×
524
                endpointSlicesSynced:          endpointSliceInformer.Informer().HasSynced,
×
525
                addOrUpdateEndpointSliceQueue: newTypedRateLimitingQueue[string]("UpdateEndpointSlice", nil),
×
526
                epKeyMutex:                    keymutex.NewHashed(numKeyLocks),
×
527

×
528
                deploymentsLister: deploymentInformer.Lister(),
×
529
                deploymentsSynced: deploymentInformer.Informer().HasSynced,
×
530

×
531
                qosPoliciesLister:    qosPolicyInformer.Lister(),
×
532
                qosPolicySynced:      qosPolicyInformer.Informer().HasSynced,
×
533
                addQoSPolicyQueue:    newTypedRateLimitingQueue("AddQoSPolicy", custCrdRateLimiter),
×
534
                updateQoSPolicyQueue: newTypedRateLimitingQueue("UpdateQoSPolicy", custCrdRateLimiter),
×
535
                delQoSPolicyQueue:    newTypedRateLimitingQueue("DeleteQoSPolicy", custCrdRateLimiter),
×
536

×
537
                configMapsLister: configMapInformer.Lister(),
×
538
                configMapsSynced: configMapInformer.Informer().HasSynced,
×
539

×
540
                sgKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
541
                sgsLister:          sgInformer.Lister(),
×
542
                sgSynced:           sgInformer.Informer().HasSynced,
×
543
                addOrUpdateSgQueue: newTypedRateLimitingQueue[string]("UpdateSecurityGroup", nil),
×
544
                delSgQueue:         newTypedRateLimitingQueue[string]("DeleteSecurityGroup", nil),
×
545
                syncSgPortsQueue:   newTypedRateLimitingQueue[string]("SyncSecurityGroupPorts", nil),
×
546

×
547
                ovnEipsLister:     ovnEipInformer.Lister(),
×
548
                ovnEipSynced:      ovnEipInformer.Informer().HasSynced,
×
549
                addOvnEipQueue:    newTypedRateLimitingQueue("AddOvnEip", custCrdRateLimiter),
×
550
                updateOvnEipQueue: newTypedRateLimitingQueue("UpdateOvnEip", custCrdRateLimiter),
×
551
                resetOvnEipQueue:  newTypedRateLimitingQueue("ResetOvnEip", custCrdRateLimiter),
×
552
                delOvnEipQueue:    newTypedRateLimitingQueue("DeleteOvnEip", custCrdRateLimiter),
×
553

×
554
                ovnFipsLister:     ovnFipInformer.Lister(),
×
555
                ovnFipSynced:      ovnFipInformer.Informer().HasSynced,
×
556
                addOvnFipQueue:    newTypedRateLimitingQueue("AddOvnFip", custCrdRateLimiter),
×
557
                updateOvnFipQueue: newTypedRateLimitingQueue("UpdateOvnFip", custCrdRateLimiter),
×
558
                delOvnFipQueue:    newTypedRateLimitingQueue("DeleteOvnFip", custCrdRateLimiter),
×
559

×
560
                ovnSnatRulesLister:     ovnSnatRuleInformer.Lister(),
×
561
                ovnSnatRuleSynced:      ovnSnatRuleInformer.Informer().HasSynced,
×
562
                addOvnSnatRuleQueue:    newTypedRateLimitingQueue("AddOvnSnatRule", custCrdRateLimiter),
×
563
                updateOvnSnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnSnatRule", custCrdRateLimiter),
×
564
                delOvnSnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnSnatRule", custCrdRateLimiter),
×
565

×
566
                ovnDnatRulesLister:     ovnDnatRuleInformer.Lister(),
×
567
                ovnDnatRuleSynced:      ovnDnatRuleInformer.Informer().HasSynced,
×
568
                addOvnDnatRuleQueue:    newTypedRateLimitingQueue("AddOvnDnatRule", custCrdRateLimiter),
×
569
                updateOvnDnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnDnatRule", custCrdRateLimiter),
×
570
                delOvnDnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnDnatRule", custCrdRateLimiter),
×
571

×
572
                csrLister:           csrInformer.Lister(),
×
573
                csrSynced:           csrInformer.Informer().HasSynced,
×
574
                addOrUpdateCsrQueue: newTypedRateLimitingQueue[string]("AddOrUpdateCSR", custCrdRateLimiter),
×
575

×
576
                addOrUpdateVMIMigrationQueue: newTypedRateLimitingQueue[string]("AddOrUpdateVMIMigration", nil),
×
577
                deleteVMQueue:                newTypedRateLimitingQueue[string]("DeleteVM", nil),
×
578
                kubevirtInformerFactory:      kubevirtInformerFactory,
×
579

×
580
                netAttachLister:          netAttachInformer.Lister(),
×
581
                netAttachSynced:          netAttachInformer.Informer().HasSynced,
×
582
                netAttachInformerFactory: attachNetInformerFactory,
×
583

×
584
                recorder:               recorder,
×
585
                informerFactory:        informerFactory,
×
586
                cmInformerFactory:      cmInformerFactory,
×
587
                deployInformerFactory:  deployInformerFactory,
×
588
                kubeovnInformerFactory: kubeovnInformerFactory,
×
589
                anpInformerFactory:     anpInformerFactory,
×
590
        }
×
591

×
592
        if controller.OVNNbClient, err = ovs.NewOvnNbClient(
×
593
                config.OvnNbAddr,
×
594
                config.OvnTimeout,
×
595
                config.OvsDbConnectTimeout,
×
596
                config.OvsDbInactivityTimeout,
×
597
                config.OvsDbConnectMaxRetry,
×
598
        ); err != nil {
×
599
                util.LogFatalAndExit(err, "failed to create ovn nb client")
×
600
        }
×
601
        if controller.OVNSbClient, err = ovs.NewOvnSbClient(
×
602
                config.OvnSbAddr,
×
603
                config.OvnTimeout,
×
604
                config.OvsDbConnectTimeout,
×
605
                config.OvsDbInactivityTimeout,
×
606
                config.OvsDbConnectMaxRetry,
×
607
        ); err != nil {
×
608
                util.LogFatalAndExit(err, "failed to create ovn sb client")
×
609
        }
×
610
        if config.EnableLb {
×
611
                controller.switchLBRuleLister = switchLBRuleInformer.Lister()
×
612
                controller.switchLBRuleSynced = switchLBRuleInformer.Informer().HasSynced
×
613
                controller.addSwitchLBRuleQueue = newTypedRateLimitingQueue("AddSwitchLBRule", custCrdRateLimiter)
×
614
                controller.delSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
615
                        "DeleteSwitchLBRule",
×
616
                        workqueue.NewTypedMaxOfRateLimiter(
×
617
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SlrInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
618
                                &workqueue.TypedBucketRateLimiter[*SlrInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
619
                        ),
×
620
                )
×
621
                controller.updateSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
622
                        "UpdateSwitchLBRule",
×
623
                        workqueue.NewTypedMaxOfRateLimiter(
×
624
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SlrInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
625
                                &workqueue.TypedBucketRateLimiter[*SlrInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
626
                        ),
×
627
                )
×
628

×
629
                controller.vpcDNSLister = vpcDNSInformer.Lister()
×
630
                controller.vpcDNSSynced = vpcDNSInformer.Informer().HasSynced
×
631
                controller.addOrUpdateVpcDNSQueue = newTypedRateLimitingQueue("AddOrUpdateVpcDns", custCrdRateLimiter)
×
632
                controller.delVpcDNSQueue = newTypedRateLimitingQueue("DeleteVpcDns", custCrdRateLimiter)
×
633
        }
×
634

635
        if config.EnableNP {
×
636
                controller.npsLister = npInformer.Lister()
×
637
                controller.npsSynced = npInformer.Informer().HasSynced
×
638
                controller.updateNpQueue = newTypedRateLimitingQueue[string]("UpdateNetworkPolicy", nil)
×
639
                controller.deleteNpQueue = newTypedRateLimitingQueue[string]("DeleteNetworkPolicy", nil)
×
640
                controller.npKeyMutex = keymutex.NewHashed(numKeyLocks)
×
641
        }
×
642

643
        if config.EnableANP {
×
644
                controller.anpsLister = anpInformer.Lister()
×
645
                controller.anpsSynced = anpInformer.Informer().HasSynced
×
646
                controller.addAnpQueue = newTypedRateLimitingQueue[string]("AddAdminNetworkPolicy", nil)
×
647
                controller.updateAnpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateAdminNetworkPolicy", nil)
×
648
                controller.deleteAnpQueue = newTypedRateLimitingQueue[*v1alpha1.AdminNetworkPolicy]("DeleteAdminNetworkPolicy", nil)
×
649
                controller.anpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
650

×
651
                controller.banpsLister = banpInformer.Lister()
×
652
                controller.banpsSynced = banpInformer.Informer().HasSynced
×
653
                controller.addBanpQueue = newTypedRateLimitingQueue[string]("AddBaseAdminNetworkPolicy", nil)
×
654
                controller.updateBanpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateBaseAdminNetworkPolicy", nil)
×
655
                controller.deleteBanpQueue = newTypedRateLimitingQueue[*v1alpha1.BaselineAdminNetworkPolicy]("DeleteBaseAdminNetworkPolicy", nil)
×
656
                controller.banpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
657
        }
×
658

659
        if config.EnableDNSNameResolver {
×
660
                controller.dnsNameResolversLister = dnsNameResolverInformer.Lister()
×
661
                controller.dnsNameResolversSynced = dnsNameResolverInformer.Informer().HasSynced
×
662
                controller.addOrUpdateDNSNameResolverQueue = newTypedRateLimitingQueue[string]("AddOrUpdateDNSNameResolver", nil)
×
663
                controller.deleteDNSNameResolverQueue = newTypedRateLimitingQueue[*kubeovnv1.DNSNameResolver]("DeleteDNSNameResolver", nil)
×
664
        }
×
665

666
        defer controller.shutdown()
×
667
        klog.Info("Starting OVN controller")
×
668

×
669
        // Wait for the caches to be synced before starting workers
×
670
        controller.informerFactory.Start(ctx.Done())
×
671
        controller.cmInformerFactory.Start(ctx.Done())
×
672
        controller.deployInformerFactory.Start(ctx.Done())
×
673
        controller.kubeovnInformerFactory.Start(ctx.Done())
×
674
        controller.anpInformerFactory.Start(ctx.Done())
×
675
        controller.StartKubevirtInformerFactory(ctx, kubevirtInformerFactory)
×
676
        controller.StartNetAttachInformerFactory(ctx)
×
677

×
678
        klog.Info("Waiting for informer caches to sync")
×
679
        cacheSyncs := []cache.InformerSynced{
×
680
                controller.vpcNatGatewaySynced, controller.vpcEgressGatewaySynced,
×
681
                controller.vpcSynced, controller.subnetSynced,
×
682
                controller.ipSynced, controller.virtualIpsSynced, controller.iptablesEipSynced,
×
683
                controller.iptablesFipSynced, controller.iptablesDnatRuleSynced, controller.iptablesSnatRuleSynced,
×
684
                controller.vlanSynced, controller.podsSynced, controller.namespacesSynced, controller.nodesSynced,
×
685
                controller.serviceSynced, controller.endpointSlicesSynced, controller.deploymentsSynced, controller.configMapsSynced,
×
686
                controller.ovnEipSynced, controller.ovnFipSynced, controller.ovnSnatRuleSynced,
×
687
                controller.ovnDnatRuleSynced,
×
688
        }
×
689
        if controller.config.EnableLb {
×
690
                cacheSyncs = append(cacheSyncs, controller.switchLBRuleSynced, controller.vpcDNSSynced)
×
691
        }
×
692
        if controller.config.EnableNP {
×
693
                cacheSyncs = append(cacheSyncs, controller.npsSynced)
×
694
        }
×
695
        if controller.config.EnableANP {
×
696
                cacheSyncs = append(cacheSyncs, controller.anpsSynced, controller.banpsSynced)
×
697
        }
×
698
        if controller.config.EnableDNSNameResolver {
×
699
                cacheSyncs = append(cacheSyncs, controller.dnsNameResolversSynced)
×
700
        }
×
701

702
        if !cache.WaitForCacheSync(ctx.Done(), cacheSyncs...) {
×
703
                util.LogFatalAndExit(nil, "failed to wait for caches to sync")
×
704
        }
×
705

706
        if _, err = podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
707
                AddFunc:    controller.enqueueAddPod,
×
708
                DeleteFunc: controller.enqueueDeletePod,
×
709
                UpdateFunc: controller.enqueueUpdatePod,
×
710
        }); err != nil {
×
711
                util.LogFatalAndExit(err, "failed to add pod event handler")
×
712
        }
×
713

714
        if _, err = namespaceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
715
                AddFunc:    controller.enqueueAddNamespace,
×
716
                UpdateFunc: controller.enqueueUpdateNamespace,
×
717
                DeleteFunc: controller.enqueueDeleteNamespace,
×
718
        }); err != nil {
×
719
                util.LogFatalAndExit(err, "failed to add namespace event handler")
×
720
        }
×
721

722
        if _, err = nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
723
                AddFunc:    controller.enqueueAddNode,
×
724
                UpdateFunc: controller.enqueueUpdateNode,
×
725
                DeleteFunc: controller.enqueueDeleteNode,
×
726
        }); err != nil {
×
727
                util.LogFatalAndExit(err, "failed to add node event handler")
×
728
        }
×
729

730
        if _, err = serviceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
731
                AddFunc:    controller.enqueueAddService,
×
732
                DeleteFunc: controller.enqueueDeleteService,
×
733
                UpdateFunc: controller.enqueueUpdateService,
×
734
        }); err != nil {
×
735
                util.LogFatalAndExit(err, "failed to add service event handler")
×
736
        }
×
737

738
        if _, err = endpointSliceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
739
                AddFunc:    controller.enqueueAddEndpointSlice,
×
740
                UpdateFunc: controller.enqueueUpdateEndpointSlice,
×
741
        }); err != nil {
×
742
                util.LogFatalAndExit(err, "failed to add endpoint slice event handler")
×
743
        }
×
744

745
        if _, err = deploymentInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
746
                AddFunc:    controller.enqueueAddDeployment,
×
747
                UpdateFunc: controller.enqueueUpdateDeployment,
×
748
        }); err != nil {
×
749
                util.LogFatalAndExit(err, "failed to add deployment event handler")
×
750
        }
×
751

752
        if _, err = vpcInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
753
                AddFunc:    controller.enqueueAddVpc,
×
754
                UpdateFunc: controller.enqueueUpdateVpc,
×
755
                DeleteFunc: controller.enqueueDelVpc,
×
756
        }); err != nil {
×
757
                util.LogFatalAndExit(err, "failed to add vpc event handler")
×
758
        }
×
759

760
        if _, err = vpcNatGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
761
                AddFunc:    controller.enqueueAddVpcNatGw,
×
762
                UpdateFunc: controller.enqueueUpdateVpcNatGw,
×
763
                DeleteFunc: controller.enqueueDeleteVpcNatGw,
×
764
        }); err != nil {
×
765
                util.LogFatalAndExit(err, "failed to add vpc nat gateway event handler")
×
766
        }
×
767

768
        if _, err = vpcEgressGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
769
                AddFunc:    controller.enqueueAddVpcEgressGateway,
×
770
                UpdateFunc: controller.enqueueUpdateVpcEgressGateway,
×
771
                DeleteFunc: controller.enqueueDeleteVpcEgressGateway,
×
772
        }); err != nil {
×
773
                util.LogFatalAndExit(err, "failed to add vpc egress gateway event handler")
×
774
        }
×
775

776
        if _, err = subnetInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
777
                AddFunc:    controller.enqueueAddSubnet,
×
778
                UpdateFunc: controller.enqueueUpdateSubnet,
×
779
                DeleteFunc: controller.enqueueDeleteSubnet,
×
780
        }); err != nil {
×
781
                util.LogFatalAndExit(err, "failed to add subnet event handler")
×
782
        }
×
783

784
        if _, err = ippoolInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
785
                AddFunc:    controller.enqueueAddIPPool,
×
786
                UpdateFunc: controller.enqueueUpdateIPPool,
×
787
                DeleteFunc: controller.enqueueDeleteIPPool,
×
788
        }); err != nil {
×
789
                util.LogFatalAndExit(err, "failed to add ippool event handler")
×
790
        }
×
791

792
        if _, err = ipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
793
                AddFunc:    controller.enqueueAddIP,
×
794
                UpdateFunc: controller.enqueueUpdateIP,
×
795
                DeleteFunc: controller.enqueueDelIP,
×
796
        }); err != nil {
×
797
                util.LogFatalAndExit(err, "failed to add ips event handler")
×
798
        }
×
799

800
        if _, err = vlanInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
801
                AddFunc:    controller.enqueueAddVlan,
×
802
                DeleteFunc: controller.enqueueDelVlan,
×
803
                UpdateFunc: controller.enqueueUpdateVlan,
×
804
        }); err != nil {
×
805
                util.LogFatalAndExit(err, "failed to add vlan event handler")
×
806
        }
×
807

808
        if _, err = sgInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
809
                AddFunc:    controller.enqueueAddSg,
×
810
                DeleteFunc: controller.enqueueDeleteSg,
×
811
                UpdateFunc: controller.enqueueUpdateSg,
×
812
        }); err != nil {
×
813
                util.LogFatalAndExit(err, "failed to add security group event handler")
×
814
        }
×
815

816
        if _, err = virtualIPInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
817
                AddFunc:    controller.enqueueAddVirtualIP,
×
818
                UpdateFunc: controller.enqueueUpdateVirtualIP,
×
819
                DeleteFunc: controller.enqueueDelVirtualIP,
×
820
        }); err != nil {
×
821
                util.LogFatalAndExit(err, "failed to add virtual ip event handler")
×
822
        }
×
823

824
        if _, err = iptablesEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
825
                AddFunc:    controller.enqueueAddIptablesEip,
×
826
                UpdateFunc: controller.enqueueUpdateIptablesEip,
×
827
                DeleteFunc: controller.enqueueDelIptablesEip,
×
828
        }); err != nil {
×
829
                util.LogFatalAndExit(err, "failed to add iptables eip event handler")
×
830
        }
×
831

832
        if _, err = iptablesFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
833
                AddFunc:    controller.enqueueAddIptablesFip,
×
834
                UpdateFunc: controller.enqueueUpdateIptablesFip,
×
835
                DeleteFunc: controller.enqueueDelIptablesFip,
×
836
        }); err != nil {
×
837
                util.LogFatalAndExit(err, "failed to add iptables fip event handler")
×
838
        }
×
839

840
        if _, err = iptablesDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
841
                AddFunc:    controller.enqueueAddIptablesDnatRule,
×
842
                UpdateFunc: controller.enqueueUpdateIptablesDnatRule,
×
843
                DeleteFunc: controller.enqueueDelIptablesDnatRule,
×
844
        }); err != nil {
×
845
                util.LogFatalAndExit(err, "failed to add iptables dnat event handler")
×
846
        }
×
847

848
        if _, err = iptablesSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
849
                AddFunc:    controller.enqueueAddIptablesSnatRule,
×
850
                UpdateFunc: controller.enqueueUpdateIptablesSnatRule,
×
851
                DeleteFunc: controller.enqueueDelIptablesSnatRule,
×
852
        }); err != nil {
×
853
                util.LogFatalAndExit(err, "failed to add iptables snat rule event handler")
×
854
        }
×
855

856
        if _, err = ovnEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
857
                AddFunc:    controller.enqueueAddOvnEip,
×
858
                UpdateFunc: controller.enqueueUpdateOvnEip,
×
859
                DeleteFunc: controller.enqueueDelOvnEip,
×
860
        }); err != nil {
×
861
                util.LogFatalAndExit(err, "failed to add ovn eip event handler")
×
862
        }
×
863

864
        if _, err = ovnFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
865
                AddFunc:    controller.enqueueAddOvnFip,
×
866
                UpdateFunc: controller.enqueueUpdateOvnFip,
×
867
                DeleteFunc: controller.enqueueDelOvnFip,
×
868
        }); err != nil {
×
869
                util.LogFatalAndExit(err, "failed to add ovn fip event handler")
×
870
        }
×
871

872
        if _, err = ovnSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
873
                AddFunc:    controller.enqueueAddOvnSnatRule,
×
874
                UpdateFunc: controller.enqueueUpdateOvnSnatRule,
×
875
                DeleteFunc: controller.enqueueDelOvnSnatRule,
×
876
        }); err != nil {
×
877
                util.LogFatalAndExit(err, "failed to add ovn snat rule event handler")
×
878
        }
×
879

880
        if _, err = ovnDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
881
                AddFunc:    controller.enqueueAddOvnDnatRule,
×
882
                UpdateFunc: controller.enqueueUpdateOvnDnatRule,
×
883
                DeleteFunc: controller.enqueueDelOvnDnatRule,
×
884
        }); err != nil {
×
885
                util.LogFatalAndExit(err, "failed to add ovn dnat rule event handler")
×
886
        }
×
887

888
        if _, err = qosPolicyInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
889
                AddFunc:    controller.enqueueAddQoSPolicy,
×
890
                UpdateFunc: controller.enqueueUpdateQoSPolicy,
×
891
                DeleteFunc: controller.enqueueDelQoSPolicy,
×
892
        }); err != nil {
×
893
                util.LogFatalAndExit(err, "failed to add qos policy event handler")
×
894
        }
×
895

896
        if config.EnableLb {
×
897
                if _, err = switchLBRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
898
                        AddFunc:    controller.enqueueAddSwitchLBRule,
×
899
                        UpdateFunc: controller.enqueueUpdateSwitchLBRule,
×
900
                        DeleteFunc: controller.enqueueDeleteSwitchLBRule,
×
901
                }); err != nil {
×
902
                        util.LogFatalAndExit(err, "failed to add switch lb rule event handler")
×
903
                }
×
904

905
                if _, err = vpcDNSInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
906
                        AddFunc:    controller.enqueueAddVpcDNS,
×
907
                        UpdateFunc: controller.enqueueUpdateVpcDNS,
×
908
                        DeleteFunc: controller.enqueueDeleteVPCDNS,
×
909
                }); err != nil {
×
910
                        util.LogFatalAndExit(err, "failed to add vpc dns event handler")
×
911
                }
×
912
        }
913

914
        if config.EnableNP {
×
915
                if _, err = npInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
916
                        AddFunc:    controller.enqueueAddNp,
×
917
                        UpdateFunc: controller.enqueueUpdateNp,
×
918
                        DeleteFunc: controller.enqueueDeleteNp,
×
919
                }); err != nil {
×
920
                        util.LogFatalAndExit(err, "failed to add network policy event handler")
×
921
                }
×
922
        }
923

924
        if config.EnableANP {
×
925
                if _, err = anpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
926
                        AddFunc:    controller.enqueueAddAnp,
×
927
                        UpdateFunc: controller.enqueueUpdateAnp,
×
928
                        DeleteFunc: controller.enqueueDeleteAnp,
×
929
                }); err != nil {
×
930
                        util.LogFatalAndExit(err, "failed to add admin network policy event handler")
×
931
                }
×
932

933
                if _, err = banpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
934
                        AddFunc:    controller.enqueueAddBanp,
×
935
                        UpdateFunc: controller.enqueueUpdateBanp,
×
936
                        DeleteFunc: controller.enqueueDeleteBanp,
×
937
                }); err != nil {
×
938
                        util.LogFatalAndExit(err, "failed to add baseline admin network policy event handler")
×
939
                }
×
940

941
                controller.anpPrioNameMap = make(map[int32]string, 100)
×
942
                controller.anpNamePrioMap = make(map[string]int32, 100)
×
943
        }
944

945
        if config.EnableDNSNameResolver {
×
946
                if _, err = dnsNameResolverInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
947
                        AddFunc:    controller.enqueueAddDNSNameResolver,
×
948
                        UpdateFunc: controller.enqueueUpdateDNSNameResolver,
×
949
                        DeleteFunc: controller.enqueueDeleteDNSNameResolver,
×
950
                }); err != nil {
×
951
                        util.LogFatalAndExit(err, "failed to add dns name resolver event handler")
×
952
                }
×
953
        }
954

955
        if config.EnableOVNIPSec {
×
956
                if _, err = csrInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
957
                        AddFunc:    controller.enqueueAddCsr,
×
958
                        UpdateFunc: controller.enqueueUpdateCsr,
×
959
                        // no need to add delete func for csr
×
960
                }); err != nil {
×
961
                        util.LogFatalAndExit(err, "failed to add csr event handler")
×
962
                }
×
963
        }
964

965
        controller.Run(ctx)
×
966
}
967

968
// Run will set up the event handlers for types we are interested in, as well
969
// as syncing informer caches and starting workers. It will block until stopCh
970
// is closed, at which point it will shutdown the workqueue and wait for
971
// workers to finish processing their current work items.
972
func (c *Controller) Run(ctx context.Context) {
×
973
        // The init process can only be placed here if the init process do really affect the normal process of controller, such as Nodes/Pods/Subnets...
×
974
        // Otherwise, the init process should be placed after all workers have already started working
×
975
        if err := c.OVNNbClient.SetLsDnatModDlDst(c.config.LsDnatModDlDst); err != nil {
×
976
                util.LogFatalAndExit(err, "failed to set NB_Global option ls_dnat_mod_dl_dst")
×
977
        }
×
978

979
        if err := c.OVNNbClient.SetUseCtInvMatch(); err != nil {
×
980
                util.LogFatalAndExit(err, "failed to set NB_Global option use_ct_inv_match to false")
×
981
        }
×
982

983
        if err := c.OVNNbClient.SetLsCtSkipDstLportIPs(c.config.LsCtSkipDstLportIPs); err != nil {
×
984
                util.LogFatalAndExit(err, "failed to set NB_Global option ls_ct_skip_dst_lport_ips")
×
985
        }
×
986

987
        if err := c.OVNNbClient.SetNodeLocalDNSIP(strings.Join(c.config.NodeLocalDNSIPs, ",")); err != nil {
×
988
                util.LogFatalAndExit(err, "failed to set NB_Global option node_local_dns_ip")
×
989
        }
×
990

991
        if err := c.OVNNbClient.SetOVNIPSec(c.config.EnableOVNIPSec); err != nil {
×
992
                util.LogFatalAndExit(err, "failed to set NB_Global ipsec")
×
993
        }
×
994

995
        if err := c.InitOVN(); err != nil {
×
996
                util.LogFatalAndExit(err, "failed to initialize ovn resources")
×
997
        }
×
998

999
        // sync ip crd before initIPAM since ip crd will be used to restore vm and statefulset pod in initIPAM
1000
        if err := c.syncIPCR(); err != nil {
×
1001
                util.LogFatalAndExit(err, "failed to sync crd ips")
×
1002
        }
×
1003

1004
        if err := c.syncFinalizers(); err != nil {
×
1005
                util.LogFatalAndExit(err, "failed to initialize crd finalizers")
×
1006
        }
×
1007

1008
        if err := c.InitIPAM(); err != nil {
×
1009
                util.LogFatalAndExit(err, "failed to initialize ipam")
×
1010
        }
×
1011

1012
        if err := c.syncNodeRoutes(); err != nil {
×
1013
                util.LogFatalAndExit(err, "failed to initialize node routes")
×
1014
        }
×
1015

1016
        if err := c.syncSubnetCR(); err != nil {
×
1017
                util.LogFatalAndExit(err, "failed to sync crd subnets")
×
1018
        }
×
1019

1020
        if err := c.syncVlanCR(); err != nil {
×
1021
                util.LogFatalAndExit(err, "failed to sync crd vlans")
×
1022
        }
×
1023

1024
        if c.config.EnableOVNIPSec && !c.config.CertManagerIPSecCert {
×
1025
                if err := c.InitDefaultOVNIPsecCA(); err != nil {
×
1026
                        util.LogFatalAndExit(err, "failed to init ovn ipsec CA")
×
1027
                }
×
1028
        }
1029

1030
        // start workers to do all the network operations
1031
        c.startWorkers(ctx)
×
1032

×
1033
        c.initResourceOnce()
×
1034
        <-ctx.Done()
×
1035
        klog.Info("Shutting down workers")
×
1036
}
1037

1038
func (c *Controller) dbStatus() {
×
1039
        const maxFailures = 5
×
1040

×
1041
        done := make(chan error, 2)
×
1042
        go func() {
×
1043
                done <- c.OVNNbClient.Echo(context.Background())
×
1044
        }()
×
1045
        go func() {
×
1046
                done <- c.OVNSbClient.Echo(context.Background())
×
1047
        }()
×
1048

1049
        resultsReceived := 0
×
1050
        timeout := time.After(time.Duration(c.config.OvnTimeout) * time.Second)
×
1051

×
1052
        for resultsReceived < 2 {
×
1053
                select {
×
1054
                case err := <-done:
×
1055
                        resultsReceived++
×
1056
                        if err != nil {
×
1057
                                c.dbFailureCount++
×
1058
                                klog.Errorf("OVN database echo failed (%d/%d): %v", c.dbFailureCount, maxFailures, err)
×
1059
                                if c.dbFailureCount >= maxFailures {
×
1060
                                        util.LogFatalAndExit(err, "OVN database connection failed after %d attempts", maxFailures)
×
1061
                                }
×
1062
                                return
×
1063
                        }
1064
                case <-timeout:
×
1065
                        c.dbFailureCount++
×
1066
                        klog.Errorf("OVN database echo timeout (%d/%d) after %ds", c.dbFailureCount, maxFailures, c.config.OvnTimeout)
×
1067
                        if c.dbFailureCount >= maxFailures {
×
1068
                                util.LogFatalAndExit(nil, "OVN database connection timeout after %d attempts", maxFailures)
×
1069
                        }
×
1070
                        return
×
1071
                }
1072
        }
1073

1074
        if c.dbFailureCount > 0 {
×
1075
                klog.Infof("OVN database connection recovered after %d failures", c.dbFailureCount)
×
1076
                c.dbFailureCount = 0
×
1077
        }
×
1078
}
1079

1080
func (c *Controller) shutdown() {
×
1081
        utilruntime.HandleCrash()
×
1082

×
1083
        c.addOrUpdatePodQueue.ShutDown()
×
1084
        c.deletePodQueue.ShutDown()
×
1085
        c.updatePodSecurityQueue.ShutDown()
×
1086

×
1087
        c.addNamespaceQueue.ShutDown()
×
1088

×
1089
        c.addOrUpdateSubnetQueue.ShutDown()
×
1090
        c.deleteSubnetQueue.ShutDown()
×
1091
        c.updateSubnetStatusQueue.ShutDown()
×
1092
        c.syncVirtualPortsQueue.ShutDown()
×
1093

×
1094
        c.addOrUpdateIPPoolQueue.ShutDown()
×
1095
        c.updateIPPoolStatusQueue.ShutDown()
×
1096
        c.deleteIPPoolQueue.ShutDown()
×
1097

×
1098
        c.addNodeQueue.ShutDown()
×
1099
        c.updateNodeQueue.ShutDown()
×
1100
        c.deleteNodeQueue.ShutDown()
×
1101

×
1102
        c.addServiceQueue.ShutDown()
×
1103
        c.deleteServiceQueue.ShutDown()
×
1104
        c.updateServiceQueue.ShutDown()
×
1105
        c.addOrUpdateEndpointSliceQueue.ShutDown()
×
1106

×
1107
        c.addVlanQueue.ShutDown()
×
1108
        c.delVlanQueue.ShutDown()
×
1109
        c.updateVlanQueue.ShutDown()
×
1110

×
1111
        c.addOrUpdateVpcQueue.ShutDown()
×
1112
        c.updateVpcStatusQueue.ShutDown()
×
1113
        c.delVpcQueue.ShutDown()
×
1114

×
1115
        c.addOrUpdateVpcNatGatewayQueue.ShutDown()
×
1116
        c.initVpcNatGatewayQueue.ShutDown()
×
1117
        c.delVpcNatGatewayQueue.ShutDown()
×
1118
        c.updateVpcEipQueue.ShutDown()
×
1119
        c.updateVpcFloatingIPQueue.ShutDown()
×
1120
        c.updateVpcDnatQueue.ShutDown()
×
1121
        c.updateVpcSnatQueue.ShutDown()
×
1122
        c.updateVpcSubnetQueue.ShutDown()
×
1123

×
1124
        c.addOrUpdateVpcEgressGatewayQueue.ShutDown()
×
1125
        c.delVpcEgressGatewayQueue.ShutDown()
×
1126

×
1127
        if c.config.EnableLb {
×
1128
                c.addSwitchLBRuleQueue.ShutDown()
×
1129
                c.delSwitchLBRuleQueue.ShutDown()
×
1130
                c.updateSwitchLBRuleQueue.ShutDown()
×
1131

×
1132
                c.addOrUpdateVpcDNSQueue.ShutDown()
×
1133
                c.delVpcDNSQueue.ShutDown()
×
1134
        }
×
1135

1136
        c.addIPQueue.ShutDown()
×
1137
        c.updateIPQueue.ShutDown()
×
1138
        c.delIPQueue.ShutDown()
×
1139

×
1140
        c.addVirtualIPQueue.ShutDown()
×
1141
        c.updateVirtualIPQueue.ShutDown()
×
1142
        c.updateVirtualParentsQueue.ShutDown()
×
1143
        c.delVirtualIPQueue.ShutDown()
×
1144

×
1145
        c.addIptablesEipQueue.ShutDown()
×
1146
        c.updateIptablesEipQueue.ShutDown()
×
1147
        c.resetIptablesEipQueue.ShutDown()
×
1148
        c.delIptablesEipQueue.ShutDown()
×
1149

×
1150
        c.addIptablesFipQueue.ShutDown()
×
1151
        c.updateIptablesFipQueue.ShutDown()
×
1152
        c.delIptablesFipQueue.ShutDown()
×
1153

×
1154
        c.addIptablesDnatRuleQueue.ShutDown()
×
1155
        c.updateIptablesDnatRuleQueue.ShutDown()
×
1156
        c.delIptablesDnatRuleQueue.ShutDown()
×
1157

×
1158
        c.addIptablesSnatRuleQueue.ShutDown()
×
1159
        c.updateIptablesSnatRuleQueue.ShutDown()
×
1160
        c.delIptablesSnatRuleQueue.ShutDown()
×
1161

×
1162
        c.addQoSPolicyQueue.ShutDown()
×
1163
        c.updateQoSPolicyQueue.ShutDown()
×
1164
        c.delQoSPolicyQueue.ShutDown()
×
1165

×
1166
        c.addOvnEipQueue.ShutDown()
×
1167
        c.updateOvnEipQueue.ShutDown()
×
1168
        c.resetOvnEipQueue.ShutDown()
×
1169
        c.delOvnEipQueue.ShutDown()
×
1170

×
1171
        c.addOvnFipQueue.ShutDown()
×
1172
        c.updateOvnFipQueue.ShutDown()
×
1173
        c.delOvnFipQueue.ShutDown()
×
1174

×
1175
        c.addOvnSnatRuleQueue.ShutDown()
×
1176
        c.updateOvnSnatRuleQueue.ShutDown()
×
1177
        c.delOvnSnatRuleQueue.ShutDown()
×
1178

×
1179
        c.addOvnDnatRuleQueue.ShutDown()
×
1180
        c.updateOvnDnatRuleQueue.ShutDown()
×
1181
        c.delOvnDnatRuleQueue.ShutDown()
×
1182

×
1183
        if c.config.EnableNP {
×
1184
                c.updateNpQueue.ShutDown()
×
1185
                c.deleteNpQueue.ShutDown()
×
1186
        }
×
1187
        if c.config.EnableANP {
×
1188
                c.addAnpQueue.ShutDown()
×
1189
                c.updateAnpQueue.ShutDown()
×
1190
                c.deleteAnpQueue.ShutDown()
×
1191

×
1192
                c.addBanpQueue.ShutDown()
×
1193
                c.updateBanpQueue.ShutDown()
×
1194
                c.deleteBanpQueue.ShutDown()
×
1195
        }
×
1196

1197
        if c.config.EnableDNSNameResolver {
×
1198
                c.addOrUpdateDNSNameResolverQueue.ShutDown()
×
1199
                c.deleteDNSNameResolverQueue.ShutDown()
×
1200
        }
×
1201

1202
        c.addOrUpdateSgQueue.ShutDown()
×
1203
        c.delSgQueue.ShutDown()
×
1204
        c.syncSgPortsQueue.ShutDown()
×
1205

×
1206
        c.addOrUpdateCsrQueue.ShutDown()
×
1207

×
1208
        if c.config.EnableLiveMigrationOptimize {
×
1209
                c.addOrUpdateVMIMigrationQueue.ShutDown()
×
1210
        }
×
1211
}
1212

1213
func (c *Controller) startWorkers(ctx context.Context) {
×
1214
        klog.Info("Starting workers")
×
1215

×
1216
        go wait.Until(runWorker("add/update vpc", c.addOrUpdateVpcQueue, c.handleAddOrUpdateVpc), time.Second, ctx.Done())
×
1217
        go wait.Until(runWorker("delete vpc", c.delVpcQueue, c.handleDelVpc), time.Second, ctx.Done())
×
1218
        go wait.Until(runWorker("update status of vpc", c.updateVpcStatusQueue, c.handleUpdateVpcStatus), time.Second, ctx.Done())
×
1219

×
1220
        go wait.Until(runWorker("add/update vpc nat gateway", c.addOrUpdateVpcNatGatewayQueue, c.handleAddOrUpdateVpcNatGw), time.Second, ctx.Done())
×
1221
        go wait.Until(runWorker("init vpc nat gateway", c.initVpcNatGatewayQueue, c.handleInitVpcNatGw), time.Second, ctx.Done())
×
1222
        go wait.Until(runWorker("delete vpc nat gateway", c.delVpcNatGatewayQueue, c.handleDelVpcNatGw), time.Second, ctx.Done())
×
1223
        go wait.Until(runWorker("add/update vpc egress gateway", c.addOrUpdateVpcEgressGatewayQueue, c.handleAddOrUpdateVpcEgressGateway), time.Second, ctx.Done())
×
1224
        go wait.Until(runWorker("delete vpc egress gateway", c.delVpcEgressGatewayQueue, c.handleDelVpcEgressGateway), time.Second, ctx.Done())
×
1225
        go wait.Until(runWorker("update fip for vpc nat gateway", c.updateVpcFloatingIPQueue, c.handleUpdateVpcFloatingIP), time.Second, ctx.Done())
×
1226
        go wait.Until(runWorker("update eip for vpc nat gateway", c.updateVpcEipQueue, c.handleUpdateVpcEip), time.Second, ctx.Done())
×
1227
        go wait.Until(runWorker("update dnat for vpc nat gateway", c.updateVpcDnatQueue, c.handleUpdateVpcDnat), time.Second, ctx.Done())
×
1228
        go wait.Until(runWorker("update snat for vpc nat gateway", c.updateVpcSnatQueue, c.handleUpdateVpcSnat), time.Second, ctx.Done())
×
1229
        go wait.Until(runWorker("update subnet route for vpc nat gateway", c.updateVpcSubnetQueue, c.handleUpdateNatGwSubnetRoute), time.Second, ctx.Done())
×
1230
        go wait.Until(runWorker("add/update csr", c.addOrUpdateCsrQueue, c.handleAddOrUpdateCsr), time.Second, ctx.Done())
×
1231
        // add default and join subnet and wait them ready
×
1232
        go wait.Until(runWorker("add/update subnet", c.addOrUpdateSubnetQueue, c.handleAddOrUpdateSubnet), time.Second, ctx.Done())
×
1233
        go wait.Until(runWorker("add/update ippool", c.addOrUpdateIPPoolQueue, c.handleAddOrUpdateIPPool), time.Second, ctx.Done())
×
1234
        go wait.Until(runWorker("add vlan", c.addVlanQueue, c.handleAddVlan), time.Second, ctx.Done())
×
1235
        go wait.Until(runWorker("add namespace", c.addNamespaceQueue, c.handleAddNamespace), time.Second, ctx.Done())
×
1236
        err := wait.PollUntilContextCancel(ctx, 3*time.Second, true, func(_ context.Context) (done bool, err error) {
×
1237
                subnets := []string{c.config.DefaultLogicalSwitch, c.config.NodeSwitch}
×
1238
                klog.Infof("wait for subnets %v ready", subnets)
×
1239

×
1240
                return c.allSubnetReady(subnets...)
×
1241
        })
×
1242
        if err != nil {
×
1243
                klog.Fatalf("wait default and join subnet ready, error: %v", err)
×
1244
        }
×
1245

1246
        go wait.Until(runWorker("add/update security group", c.addOrUpdateSgQueue, func(key string) error { return c.handleAddOrUpdateSg(key, false) }), time.Second, ctx.Done())
×
1247
        go wait.Until(runWorker("delete security group", c.delSgQueue, c.handleDeleteSg), time.Second, ctx.Done())
×
1248
        go wait.Until(runWorker("ports for security group", c.syncSgPortsQueue, c.syncSgLogicalPort), time.Second, ctx.Done())
×
1249

×
1250
        // run node worker before handle any pods
×
1251
        for range c.config.WorkerNum {
×
1252
                go wait.Until(runWorker("add node", c.addNodeQueue, c.handleAddNode), time.Second, ctx.Done())
×
1253
                go wait.Until(runWorker("update node", c.updateNodeQueue, c.handleUpdateNode), time.Second, ctx.Done())
×
1254
                go wait.Until(runWorker("delete node", c.deleteNodeQueue, c.handleDeleteNode), time.Second, ctx.Done())
×
1255
        }
×
1256
        for {
×
1257
                ready := true
×
1258
                time.Sleep(3 * time.Second)
×
1259
                nodes, err := c.nodesLister.List(labels.Everything())
×
1260
                if err != nil {
×
1261
                        util.LogFatalAndExit(err, "failed to list nodes")
×
1262
                }
×
1263
                for _, node := range nodes {
×
1264
                        if node.Annotations[util.AllocatedAnnotation] != "true" {
×
1265
                                klog.Infof("wait node %s annotation ready", node.Name)
×
1266
                                ready = false
×
1267
                                break
×
1268
                        }
1269
                }
1270
                if ready {
×
1271
                        break
×
1272
                }
1273
        }
1274

1275
        if c.config.EnableLb {
×
1276
                go wait.Until(runWorker("add service", c.addServiceQueue, c.handleAddService), time.Second, ctx.Done())
×
1277
                // run in a single worker to avoid delete the last vip, which will lead ovn to delete the loadbalancer
×
1278
                go wait.Until(runWorker("delete service", c.deleteServiceQueue, c.handleDeleteService), time.Second, ctx.Done())
×
1279

×
1280
                go wait.Until(runWorker("add/update switch lb rule", c.addSwitchLBRuleQueue, c.handleAddOrUpdateSwitchLBRule), time.Second, ctx.Done())
×
1281
                go wait.Until(runWorker("delete switch lb rule", c.delSwitchLBRuleQueue, c.handleDelSwitchLBRule), time.Second, ctx.Done())
×
1282
                go wait.Until(runWorker("delete switch lb rule", c.updateSwitchLBRuleQueue, c.handleUpdateSwitchLBRule), time.Second, ctx.Done())
×
1283

×
1284
                go wait.Until(runWorker("add/update vpc dns", c.addOrUpdateVpcDNSQueue, c.handleAddOrUpdateVPCDNS), time.Second, ctx.Done())
×
1285
                go wait.Until(runWorker("delete vpc dns", c.delVpcDNSQueue, c.handleDelVpcDNS), time.Second, ctx.Done())
×
1286
                go wait.Until(func() {
×
1287
                        c.resyncVpcDNSConfig()
×
1288
                }, 5*time.Second, ctx.Done())
×
1289
        }
1290

1291
        for range c.config.WorkerNum {
×
1292
                go wait.Until(runWorker("delete pod", c.deletePodQueue, c.handleDeletePod), time.Second, ctx.Done())
×
1293
                go wait.Until(runWorker("add/update pod", c.addOrUpdatePodQueue, c.handleAddOrUpdatePod), time.Second, ctx.Done())
×
1294
                go wait.Until(runWorker("update pod security", c.updatePodSecurityQueue, c.handleUpdatePodSecurity), time.Second, ctx.Done())
×
1295

×
1296
                go wait.Until(runWorker("delete subnet", c.deleteSubnetQueue, c.handleDeleteSubnet), time.Second, ctx.Done())
×
1297
                go wait.Until(runWorker("delete ippool", c.deleteIPPoolQueue, c.handleDeleteIPPool), time.Second, ctx.Done())
×
1298
                go wait.Until(runWorker("update status of subnet", c.updateSubnetStatusQueue, c.handleUpdateSubnetStatus), time.Second, ctx.Done())
×
1299
                go wait.Until(runWorker("update status of ippool", c.updateIPPoolStatusQueue, c.handleUpdateIPPoolStatus), time.Second, ctx.Done())
×
1300
                go wait.Until(runWorker("virtual port for subnet", c.syncVirtualPortsQueue, c.syncVirtualPort), time.Second, ctx.Done())
×
1301

×
1302
                if c.config.EnableLb {
×
1303
                        go wait.Until(runWorker("update service", c.updateServiceQueue, c.handleUpdateService), time.Second, ctx.Done())
×
1304
                        go wait.Until(runWorker("add/update endpoint slice", c.addOrUpdateEndpointSliceQueue, c.handleUpdateEndpointSlice), time.Second, ctx.Done())
×
1305
                }
×
1306

1307
                if c.config.EnableNP {
×
1308
                        go wait.Until(runWorker("update network policy", c.updateNpQueue, c.handleUpdateNp), time.Second, ctx.Done())
×
1309
                        go wait.Until(runWorker("delete network policy", c.deleteNpQueue, c.handleDeleteNp), time.Second, ctx.Done())
×
1310
                }
×
1311

1312
                go wait.Until(runWorker("delete vlan", c.delVlanQueue, c.handleDelVlan), time.Second, ctx.Done())
×
1313
                go wait.Until(runWorker("update vlan", c.updateVlanQueue, c.handleUpdateVlan), time.Second, ctx.Done())
×
1314
        }
1315

1316
        if c.config.EnableEipSnat {
×
1317
                go wait.Until(func() {
×
1318
                        // init l3 about the default vpc external lrp binding to the gw chassis
×
1319
                        c.resyncExternalGateway()
×
1320
                }, time.Second, ctx.Done())
×
1321

1322
                // maintain l3 ha about the vpc external lrp binding to the gw chassis
1323
                c.OVNNbClient.MonitorBFD()
×
1324
        }
1325
        // TODO: we should merge these two vpc nat config into one config and resync them together
1326
        go wait.Until(func() {
×
1327
                c.resyncVpcNatGwConfig()
×
1328
        }, time.Second, ctx.Done())
×
1329

1330
        go wait.Until(func() {
×
1331
                c.resyncVpcNatConfig()
×
1332
        }, time.Second, ctx.Done())
×
1333

1334
        if c.config.GCInterval != 0 {
×
1335
                go wait.Until(func() {
×
1336
                        if err := c.markAndCleanLSP(); err != nil {
×
1337
                                klog.Errorf("gc lsp error: %v", err)
×
1338
                        }
×
1339
                }, time.Duration(c.config.GCInterval)*time.Second, ctx.Done())
1340
        }
1341

1342
        go wait.Until(func() {
×
1343
                if err := c.inspectPod(); err != nil {
×
1344
                        klog.Errorf("inspection error: %v", err)
×
1345
                }
×
1346
        }, time.Duration(c.config.InspectInterval)*time.Second, ctx.Done())
1347

1348
        if c.config.EnableExternalVpc {
×
1349
                go wait.Until(func() {
×
1350
                        c.syncExternalVpc()
×
1351
                }, 5*time.Second, ctx.Done())
×
1352
        }
1353

1354
        go wait.Until(c.resyncProviderNetworkStatus, 30*time.Second, ctx.Done())
×
1355
        go wait.Until(c.exportSubnetMetrics, 30*time.Second, ctx.Done())
×
1356
        go wait.Until(c.checkSubnetGateway, 5*time.Second, ctx.Done())
×
1357

×
1358
        go wait.Until(runWorker("add ovn eip", c.addOvnEipQueue, c.handleAddOvnEip), time.Second, ctx.Done())
×
1359
        go wait.Until(runWorker("update ovn eip", c.updateOvnEipQueue, c.handleUpdateOvnEip), time.Second, ctx.Done())
×
1360
        go wait.Until(runWorker("reset ovn eip", c.resetOvnEipQueue, c.handleResetOvnEip), time.Second, ctx.Done())
×
1361
        go wait.Until(runWorker("delete ovn eip", c.delOvnEipQueue, c.handleDelOvnEip), time.Second, ctx.Done())
×
1362

×
1363
        go wait.Until(runWorker("add ovn fip", c.addOvnFipQueue, c.handleAddOvnFip), time.Second, ctx.Done())
×
1364
        go wait.Until(runWorker("update ovn fip", c.updateOvnFipQueue, c.handleUpdateOvnFip), time.Second, ctx.Done())
×
1365
        go wait.Until(runWorker("delete ovn fip", c.delOvnFipQueue, c.handleDelOvnFip), time.Second, ctx.Done())
×
1366

×
1367
        go wait.Until(runWorker("add ovn snat rule", c.addOvnSnatRuleQueue, c.handleAddOvnSnatRule), time.Second, ctx.Done())
×
1368
        go wait.Until(runWorker("update ovn snat rule", c.updateOvnSnatRuleQueue, c.handleUpdateOvnSnatRule), time.Second, ctx.Done())
×
1369
        go wait.Until(runWorker("delete ovn snat rule", c.delOvnSnatRuleQueue, c.handleDelOvnSnatRule), time.Second, ctx.Done())
×
1370

×
1371
        go wait.Until(runWorker("add ovn dnat", c.addOvnDnatRuleQueue, c.handleAddOvnDnatRule), time.Second, ctx.Done())
×
1372
        go wait.Until(runWorker("update ovn dnat", c.updateOvnDnatRuleQueue, c.handleUpdateOvnDnatRule), time.Second, ctx.Done())
×
1373
        go wait.Until(runWorker("delete ovn dnat", c.delOvnDnatRuleQueue, c.handleDelOvnDnatRule), time.Second, ctx.Done())
×
1374

×
1375
        go wait.Until(c.CheckNodePortGroup, time.Duration(c.config.NodePgProbeTime)*time.Minute, ctx.Done())
×
1376

×
1377
        go wait.Until(runWorker("add ip", c.addIPQueue, c.handleAddReservedIP), time.Second, ctx.Done())
×
1378
        go wait.Until(runWorker("update ip", c.updateIPQueue, c.handleUpdateIP), time.Second, ctx.Done())
×
1379
        go wait.Until(runWorker("delete ip", c.delIPQueue, c.handleDelIP), time.Second, ctx.Done())
×
1380

×
1381
        go wait.Until(runWorker("add vip", c.addVirtualIPQueue, c.handleAddVirtualIP), time.Second, ctx.Done())
×
1382
        go wait.Until(runWorker("update vip", c.updateVirtualIPQueue, c.handleUpdateVirtualIP), time.Second, ctx.Done())
×
1383
        go wait.Until(runWorker("update virtual parent for vip", c.updateVirtualParentsQueue, c.handleUpdateVirtualParents), time.Second, ctx.Done())
×
1384
        go wait.Until(runWorker("delete vip", c.delVirtualIPQueue, c.handleDelVirtualIP), time.Second, ctx.Done())
×
1385

×
1386
        go wait.Until(runWorker("add iptables eip", c.addIptablesEipQueue, c.handleAddIptablesEip), time.Second, ctx.Done())
×
1387
        go wait.Until(runWorker("update iptables eip", c.updateIptablesEipQueue, c.handleUpdateIptablesEip), time.Second, ctx.Done())
×
1388
        go wait.Until(runWorker("reset iptables eip", c.resetIptablesEipQueue, c.handleResetIptablesEip), time.Second, ctx.Done())
×
1389
        go wait.Until(runWorker("delete iptables eip", c.delIptablesEipQueue, c.handleDelIptablesEip), time.Second, ctx.Done())
×
1390

×
1391
        go wait.Until(runWorker("add iptables fip", c.addIptablesFipQueue, c.handleAddIptablesFip), time.Second, ctx.Done())
×
1392
        go wait.Until(runWorker("update iptables fip", c.updateIptablesFipQueue, c.handleUpdateIptablesFip), time.Second, ctx.Done())
×
1393
        go wait.Until(runWorker("delete iptables fip", c.delIptablesFipQueue, c.handleDelIptablesFip), time.Second, ctx.Done())
×
1394

×
1395
        go wait.Until(runWorker("add iptables dnat rule", c.addIptablesDnatRuleQueue, c.handleAddIptablesDnatRule), time.Second, ctx.Done())
×
1396
        go wait.Until(runWorker("update iptables dnat rule", c.updateIptablesDnatRuleQueue, c.handleUpdateIptablesDnatRule), time.Second, ctx.Done())
×
1397
        go wait.Until(runWorker("delete iptables dnat rule", c.delIptablesDnatRuleQueue, c.handleDelIptablesDnatRule), time.Second, ctx.Done())
×
1398

×
1399
        go wait.Until(runWorker("add iptables snat rule", c.addIptablesSnatRuleQueue, c.handleAddIptablesSnatRule), time.Second, ctx.Done())
×
1400
        go wait.Until(runWorker("update iptables snat rule", c.updateIptablesSnatRuleQueue, c.handleUpdateIptablesSnatRule), time.Second, ctx.Done())
×
1401
        go wait.Until(runWorker("delete iptables snat rule", c.delIptablesSnatRuleQueue, c.handleDelIptablesSnatRule), time.Second, ctx.Done())
×
1402

×
1403
        go wait.Until(runWorker("add qos policy", c.addQoSPolicyQueue, c.handleAddQoSPolicy), time.Second, ctx.Done())
×
1404
        go wait.Until(runWorker("update qos policy", c.updateQoSPolicyQueue, c.handleUpdateQoSPolicy), time.Second, ctx.Done())
×
1405
        go wait.Until(runWorker("delete qos policy", c.delQoSPolicyQueue, c.handleDelQoSPolicy), time.Second, ctx.Done())
×
1406

×
1407
        if c.config.EnableANP {
×
1408
                go wait.Until(runWorker("add admin network policy", c.addAnpQueue, c.handleAddAnp), time.Second, ctx.Done())
×
1409
                go wait.Until(runWorker("update admin network policy", c.updateAnpQueue, c.handleUpdateAnp), time.Second, ctx.Done())
×
1410
                go wait.Until(runWorker("delete admin network policy", c.deleteAnpQueue, c.handleDeleteAnp), time.Second, ctx.Done())
×
1411

×
1412
                go wait.Until(runWorker("add base admin network policy", c.addBanpQueue, c.handleAddBanp), time.Second, ctx.Done())
×
1413
                go wait.Until(runWorker("update base admin network policy", c.updateBanpQueue, c.handleUpdateBanp), time.Second, ctx.Done())
×
1414
                go wait.Until(runWorker("delete base admin network policy", c.deleteBanpQueue, c.handleDeleteBanp), time.Second, ctx.Done())
×
1415
        }
×
1416

1417
        if c.config.EnableDNSNameResolver {
×
1418
                go wait.Until(runWorker("add or update dns name resolver", c.addOrUpdateDNSNameResolverQueue, c.handleAddOrUpdateDNSNameResolver), time.Second, ctx.Done())
×
1419
                go wait.Until(runWorker("delete dns name resolver", c.deleteDNSNameResolverQueue, c.handleDeleteDNSNameResolver), time.Second, ctx.Done())
×
1420
        }
×
1421

1422
        if c.config.EnableLiveMigrationOptimize {
×
1423
                go wait.Until(runWorker("add/update vmiMigration ", c.addOrUpdateVMIMigrationQueue, c.handleAddOrUpdateVMIMigration), 50*time.Millisecond, ctx.Done())
×
1424
        }
×
1425

1426
        go wait.Until(runWorker("delete vm", c.deleteVMQueue, c.handleDeleteVM), time.Second, ctx.Done())
×
1427

×
1428
        go wait.Until(c.dbStatus, 15*time.Second, ctx.Done())
×
1429
}
1430

1431
func (c *Controller) allSubnetReady(subnets ...string) (bool, error) {
1✔
1432
        for _, lsName := range subnets {
2✔
1433
                exist, err := c.OVNNbClient.LogicalSwitchExists(lsName)
1✔
1434
                if err != nil {
1✔
1435
                        klog.Error(err)
×
1436
                        return false, fmt.Errorf("check logical switch %s exist: %w", lsName, err)
×
1437
                }
×
1438

1439
                if !exist {
2✔
1440
                        return false, nil
1✔
1441
                }
1✔
1442
        }
1443

1444
        return true, nil
1✔
1445
}
1446

1447
func (c *Controller) initResourceOnce() {
×
1448
        c.registerSubnetMetrics()
×
1449

×
1450
        if err := c.initNodeChassis(); err != nil {
×
1451
                util.LogFatalAndExit(err, "failed to initialize node chassis")
×
1452
        }
×
1453

1454
        if err := c.initDefaultDenyAllSecurityGroup(); err != nil {
×
1455
                util.LogFatalAndExit(err, "failed to initialize 'deny_all' security group")
×
1456
        }
×
1457
        if err := c.syncSecurityGroup(); err != nil {
×
1458
                util.LogFatalAndExit(err, "failed to sync security group")
×
1459
        }
×
1460

1461
        if err := c.syncVpcNatGatewayCR(); err != nil {
×
1462
                util.LogFatalAndExit(err, "failed to sync crd vpc nat gateways")
×
1463
        }
×
1464

1465
        if err := c.initVpcNatGw(); err != nil {
×
1466
                util.LogFatalAndExit(err, "failed to initialize vpc nat gateways")
×
1467
        }
×
1468
        if c.config.EnableLb {
×
1469
                if err := c.initVpcDNSConfig(); err != nil {
×
1470
                        util.LogFatalAndExit(err, "failed to initialize vpc-dns")
×
1471
                }
×
1472
        }
1473

1474
        // remove resources in ovndb that not exist any more in kubernetes resources
1475
        // process gc at last in case of affecting other init process
1476
        if err := c.gc(); err != nil {
×
1477
                util.LogFatalAndExit(err, "failed to run gc")
×
1478
        }
×
1479
}
1480

1481
func processNextWorkItem[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error, getItemKey func(any) string) bool {
×
1482
        item, shutdown := queue.Get()
×
1483
        if shutdown {
×
1484
                return false
×
1485
        }
×
1486

1487
        err := func(item T) error {
×
1488
                defer queue.Done(item)
×
1489
                if err := handler(item); err != nil {
×
1490
                        queue.AddRateLimited(item)
×
1491
                        return fmt.Errorf("error syncing %s %q: %w, requeuing", action, getItemKey(item), err)
×
1492
                }
×
1493
                queue.Forget(item)
×
1494
                return nil
×
1495
        }(item)
1496
        if err != nil {
×
1497
                utilruntime.HandleError(err)
×
1498
                return true
×
1499
        }
×
1500
        return true
×
1501
}
1502

1503
func getWorkItemKey(obj any) string {
×
1504
        switch v := obj.(type) {
×
1505
        case string:
×
1506
                return v
×
1507
        case *vpcService:
×
1508
                return cache.MetaObjectToName(obj.(*vpcService).Svc).String()
×
1509
        case *AdminNetworkPolicyChangedDelta:
×
1510
                return v.key
×
1511
        case *SlrInfo:
×
1512
                return v.Name
×
1513
        default:
×
1514
                key, err := cache.MetaNamespaceKeyFunc(obj)
×
1515
                if err != nil {
×
1516
                        utilruntime.HandleError(err)
×
1517
                        return ""
×
1518
                }
×
1519
                return key
×
1520
        }
1521
}
1522

1523
func runWorker[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error) func() {
×
1524
        return func() {
×
1525
                for processNextWorkItem(action, queue, handler, getWorkItemKey) {
×
1526
                }
×
1527
        }
1528
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc