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

kubeovn / kube-ovn / 29887791413

22 Jul 2026 03:10AM UTC coverage: 27.912% (+0.2%) from 27.746%
29887791413

Pull #7040

github

changluyi
fix(metallb): simplify local external VIP handling
Pull Request #7040: feature(metallb): internal underlay pod traffic to LB service VIP node

115 of 264 new or added lines in 5 files covered. (43.56%)

489 existing lines in 5 files now uncovered.

17076 of 61179 relevant lines covered (27.91%)

0.33 hits per line

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

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

3
import (
4
        "context"
5
        "fmt"
6
        "runtime"
7
        "strings"
8
        "sync"
9
        "sync/atomic"
10
        "time"
11

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

41
        kubeovnv1 "github.com/kubeovn/kube-ovn/pkg/apis/kubeovn/v1"
42
        kubeovninformer "github.com/kubeovn/kube-ovn/pkg/client/informers/externalversions"
43
        kubeovnlister "github.com/kubeovn/kube-ovn/pkg/client/listers/kubeovn/v1"
44
        "github.com/kubeovn/kube-ovn/pkg/informer"
45
        ovnipam "github.com/kubeovn/kube-ovn/pkg/ipam"
46
        "github.com/kubeovn/kube-ovn/pkg/ovs"
47
        "github.com/kubeovn/kube-ovn/pkg/util"
48
)
49

50
const controllerAgentName = "kube-ovn-controller"
51

52
const (
53
        logicalSwitchKey              = "ls"
54
        logicalRouterKey              = "lr"
55
        portGroupKey                  = "pg"
56
        networkPolicyKey              = "np"
57
        sgKey                         = "sg"
58
        sgsKey                        = "security_groups"
59
        u2oKey                        = "u2o"
60
        adminNetworkPolicyKey         = "anp"
61
        baselineAdminNetworkPolicyKey = "banp"
62
        ippoolKey                     = "ippool"
63
        clusterNetworkPolicyKey       = "cnp"
64
)
65

66
// Controller is kube-ovn main controller that watch ns/pod/node/svc/ep and operate ovn
67
type Controller struct {
68
        config *Configuration
69

70
        ipam             *ovnipam.IPAM
71
        namedPort        *NamedPort
72
        anpPrioNameMap   map[int32]string
73
        anpNamePrioMap   map[string]int32
74
        bnpPrioNameMap   map[int32]string
75
        bnpNamePrioMap   map[string]int32
76
        priorityMapMutex sync.RWMutex
77

78
        OVNNbClient ovs.NbClient
79
        OVNSbClient ovs.SbClient
80

81
        // ExternalGatewayType define external gateway type, centralized
82
        ExternalGatewayType string
83

84
        podsLister             v1.PodLister
85
        podsSynced             cache.InformerSynced
86
        podIndexer             cache.Indexer
87
        addOrUpdatePodQueue    workqueue.TypedRateLimitingInterface[string]
88
        deletePodQueue         workqueue.TypedRateLimitingInterface[string]
89
        deletingPodObjMap      *xsync.Map[string, *corev1.Pod]
90
        deletingNodeObjMap     *xsync.Map[string, *corev1.Node]
91
        updatePodSecurityQueue workqueue.TypedRateLimitingInterface[string]
92
        podKeyMutex            keymutex.KeyMutex
93

94
        vpcsLister           kubeovnlister.VpcLister
95
        vpcSynced            cache.InformerSynced
96
        vpcIndexer           cache.Indexer
97
        addOrUpdateVpcQueue  workqueue.TypedRateLimitingInterface[string]
98
        vpcLastPoliciesMap   *xsync.Map[string, string]
99
        delVpcQueue          workqueue.TypedRateLimitingInterface[*kubeovnv1.Vpc]
100
        updateVpcStatusQueue workqueue.TypedRateLimitingInterface[string]
101
        vpcKeyMutex          keymutex.KeyMutex
102

103
        vpcNatGatewayLister           kubeovnlister.VpcNatGatewayLister
104
        vpcNatGatewaySynced           cache.InformerSynced
105
        addOrUpdateVpcNatGatewayQueue workqueue.TypedRateLimitingInterface[string]
106
        delVpcNatGatewayQueue         workqueue.TypedRateLimitingInterface[string]
107
        initVpcNatGatewayQueue        workqueue.TypedRateLimitingInterface[string]
108
        updateVpcEipQueue             workqueue.TypedRateLimitingInterface[string]
109
        updateVpcFloatingIPQueue      workqueue.TypedRateLimitingInterface[string]
110
        updateVpcDnatQueue            workqueue.TypedRateLimitingInterface[string]
111
        updateVpcSnatQueue            workqueue.TypedRateLimitingInterface[string]
112
        updateVpcSubnetQueue          workqueue.TypedRateLimitingInterface[string]
113
        vpcNatGwKeyMutex              keymutex.KeyMutex
114
        vpcNatGwExecKeyMutex          keymutex.KeyMutex
115

116
        vpcEgressGatewayLister           kubeovnlister.VpcEgressGatewayLister
117
        vpcEgressGatewaySynced           cache.InformerSynced
118
        addOrUpdateVpcEgressGatewayQueue workqueue.TypedRateLimitingInterface[string]
119
        delVpcEgressGatewayQueue         workqueue.TypedRateLimitingInterface[string]
120
        vpcEgressGatewayKeyMutex         keymutex.KeyMutex
121

122
        // bgpConfLister/evpnConfLister are published asynchronously by the
123
        // optional-CRD background poller (StartBgpEvpnConfInformerFactory), but
124
        // read by VEG worker goroutines. Atomic pointers keep the read/write
125
        // race-free without locking the hot reconcile path.
126
        bgpConfLister  atomic.Pointer[kubeovnlister.BgpConfLister]
127
        bgpConfSynced  cache.InformerSynced
128
        evpnConfLister atomic.Pointer[kubeovnlister.EvpnConfLister]
129
        evpnConfSynced cache.InformerSynced
130

131
        routerLBRuleLister      kubeovnlister.RouterLBRuleLister
132
        routerLBRuleSynced      cache.InformerSynced
133
        addRouterLBRuleQueue    workqueue.TypedRateLimitingInterface[string]
134
        updateRouterLBRuleQueue workqueue.TypedRateLimitingInterface[*RouterLBRuleInfo]
135
        delRouterLBRuleQueue    workqueue.TypedRateLimitingInterface[*RouterLBRuleInfo]
136

137
        switchLBRuleLister      kubeovnlister.SwitchLBRuleLister
138
        switchLBRuleSynced      cache.InformerSynced
139
        addSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[string]
140
        updateSwitchLBRuleQueue workqueue.TypedRateLimitingInterface[*SwitchLBRuleInfo]
141
        delSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[*SwitchLBRuleInfo]
142

143
        vpcDNSLister           kubeovnlister.VpcDnsLister
144
        vpcDNSSynced           cache.InformerSynced
145
        addOrUpdateVpcDNSQueue workqueue.TypedRateLimitingInterface[string]
146
        delVpcDNSQueue         workqueue.TypedRateLimitingInterface[string]
147

148
        subnetsLister           kubeovnlister.SubnetLister
149
        subnetSynced            cache.InformerSynced
150
        addOrUpdateSubnetQueue  workqueue.TypedRateLimitingInterface[string]
151
        deleteSubnetQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.Subnet]
152
        updateSubnetStatusQueue workqueue.TypedRateLimitingInterface[string]
153
        syncVirtualPortsQueue   workqueue.TypedRateLimitingInterface[string]
154
        subnetKeyMutex          keymutex.KeyMutex
155

156
        ippoolLister            kubeovnlister.IPPoolLister
157
        ippoolSynced            cache.InformerSynced
158
        addOrUpdateIPPoolQueue  workqueue.TypedRateLimitingInterface[string]
159
        updateIPPoolStatusQueue workqueue.TypedRateLimitingInterface[string]
160
        deleteIPPoolQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.IPPool]
161
        ippoolKeyMutex          keymutex.KeyMutex
162

163
        ipsLister     kubeovnlister.IPLister
164
        ipSynced      cache.InformerSynced
165
        ipIndexer     cache.Indexer
166
        addIPQueue    workqueue.TypedRateLimitingInterface[string]
167
        updateIPQueue workqueue.TypedRateLimitingInterface[string]
168
        delIPQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.IP]
169

170
        virtualIpsLister          kubeovnlister.VipLister
171
        virtualIpsSynced          cache.InformerSynced
172
        addVirtualIPQueue         workqueue.TypedRateLimitingInterface[string]
173
        updateVirtualIPQueue      workqueue.TypedRateLimitingInterface[string]
174
        updateVirtualParentsQueue workqueue.TypedRateLimitingInterface[string]
175
        delVirtualIPQueue         workqueue.TypedRateLimitingInterface[*kubeovnv1.Vip]
176

177
        iptablesEipsLister     kubeovnlister.IptablesEIPLister
178
        iptablesEipSynced      cache.InformerSynced
179
        addIptablesEipQueue    workqueue.TypedRateLimitingInterface[string]
180
        updateIptablesEipQueue workqueue.TypedRateLimitingInterface[string]
181
        resetIptablesEipQueue  workqueue.TypedRateLimitingInterface[string]
182
        delIptablesEipQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.IptablesEIP]
183

184
        iptablesFipsLister     kubeovnlister.IptablesFIPRuleLister
185
        iptablesFipSynced      cache.InformerSynced
186
        addIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
187
        updateIptablesFipQueue workqueue.TypedRateLimitingInterface[string]
188
        delIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
189

190
        iptablesDnatRulesLister     kubeovnlister.IptablesDnatRuleLister
191
        iptablesDnatRuleSynced      cache.InformerSynced
192
        addIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
193
        updateIptablesDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
194
        delIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
195

196
        iptablesSnatRulesLister     kubeovnlister.IptablesSnatRuleLister
197
        iptablesSnatRuleSynced      cache.InformerSynced
198
        addIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
199
        updateIptablesSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
200
        delIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
201

202
        ovnEipsLister     kubeovnlister.OvnEipLister
203
        ovnEipSynced      cache.InformerSynced
204
        addOvnEipQueue    workqueue.TypedRateLimitingInterface[string]
205
        updateOvnEipQueue workqueue.TypedRateLimitingInterface[string]
206
        resetOvnEipQueue  workqueue.TypedRateLimitingInterface[string]
207
        delOvnEipQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.OvnEip]
208

209
        ovnFipsLister     kubeovnlister.OvnFipLister
210
        ovnFipSynced      cache.InformerSynced
211
        addOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
212
        updateOvnFipQueue workqueue.TypedRateLimitingInterface[string]
213
        delOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
214

215
        ovnSnatRulesLister     kubeovnlister.OvnSnatRuleLister
216
        ovnSnatRuleSynced      cache.InformerSynced
217
        addOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
218
        updateOvnSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
219
        delOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
220

221
        ovnDnatRulesLister     kubeovnlister.OvnDnatRuleLister
222
        ovnDnatRuleSynced      cache.InformerSynced
223
        addOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
224
        updateOvnDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
225
        delOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
226

227
        providerNetworksLister kubeovnlister.ProviderNetworkLister
228
        providerNetworkSynced  cache.InformerSynced
229

230
        vlansLister     kubeovnlister.VlanLister
231
        vlanSynced      cache.InformerSynced
232
        addVlanQueue    workqueue.TypedRateLimitingInterface[string]
233
        delVlanQueue    workqueue.TypedRateLimitingInterface[string]
234
        updateVlanQueue workqueue.TypedRateLimitingInterface[string]
235
        vlanKeyMutex    keymutex.KeyMutex
236

237
        namespacesLister  v1.NamespaceLister
238
        namespacesSynced  cache.InformerSynced
239
        addNamespaceQueue workqueue.TypedRateLimitingInterface[string]
240
        nsKeyMutex        keymutex.KeyMutex
241

242
        nodesLister     v1.NodeLister
243
        nodesSynced     cache.InformerSynced
244
        addNodeQueue    workqueue.TypedRateLimitingInterface[string]
245
        updateNodeQueue workqueue.TypedRateLimitingInterface[string]
246
        deleteNodeQueue workqueue.TypedRateLimitingInterface[string]
247
        nodeKeyMutex    keymutex.KeyMutex
248

249
        servicesLister     v1.ServiceLister
250
        serviceSynced      cache.InformerSynced
251
        addServiceQueue    workqueue.TypedRateLimitingInterface[string]
252
        deleteServiceQueue workqueue.TypedRateLimitingInterface[*vpcService]
253
        updateServiceQueue workqueue.TypedRateLimitingInterface[*updateSvcObject]
254
        svcKeyMutex        keymutex.KeyMutex
255

256
        endpointSlicesLister          discoveryv1.EndpointSliceLister
257
        endpointSlicesSynced          cache.InformerSynced
258
        epsIndexer                    cache.Indexer
259
        addOrUpdateEndpointSliceQueue workqueue.TypedRateLimitingInterface[string]
260
        epKeyMutex                    keymutex.KeyMutex
261
        serviceL2StatusMutex          sync.RWMutex
262
        serviceL2StatusIndexer        cache.Indexer
263
        serviceL2StatusSynced         cache.InformerSynced
264
        serviceL2StatusStarted        bool
265

266
        deploymentsLister appsv1.DeploymentLister
267
        deploymentsSynced cache.InformerSynced
268

269
        npsLister     netv1.NetworkPolicyLister
270
        npsSynced     cache.InformerSynced
271
        npIndexer     cache.Indexer
272
        updateNpQueue workqueue.TypedRateLimitingInterface[string]
273
        deleteNpQueue workqueue.TypedRateLimitingInterface[string]
274
        npKeyMutex    keymutex.KeyMutex
275

276
        sgsLister          kubeovnlister.SecurityGroupLister
277
        sgSynced           cache.InformerSynced
278
        addOrUpdateSgQueue workqueue.TypedRateLimitingInterface[string]
279
        delSgQueue         workqueue.TypedRateLimitingInterface[string]
280
        syncSgPortsQueue   workqueue.TypedRateLimitingInterface[string]
281
        sgKeyMutex         keymutex.KeyMutex
282

283
        qosPoliciesLister    kubeovnlister.QoSPolicyLister
284
        qosPolicySynced      cache.InformerSynced
285
        addQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
286
        updateQoSPolicyQueue workqueue.TypedRateLimitingInterface[string]
287
        delQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
288

289
        configMapsLister v1.ConfigMapLister
290
        configMapsSynced cache.InformerSynced
291

292
        anpsLister     anplister.AdminNetworkPolicyLister
293
        anpsSynced     cache.InformerSynced
294
        addAnpQueue    workqueue.TypedRateLimitingInterface[string]
295
        updateAnpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
296
        deleteAnpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.AdminNetworkPolicy]
297
        anpKeyMutex    keymutex.KeyMutex
298

299
        dnsNameResolversLister          kubeovnlister.DNSNameResolverLister
300
        dnsNameResolverIndexer          cache.Indexer
301
        dnsNameResolversSynced          cache.InformerSynced
302
        addOrUpdateDNSNameResolverQueue workqueue.TypedRateLimitingInterface[string]
303
        deleteDNSNameResolverQueue      workqueue.TypedRateLimitingInterface[*kubeovnv1.DNSNameResolver]
304

305
        banpsLister     anplister.BaselineAdminNetworkPolicyLister
306
        banpsSynced     cache.InformerSynced
307
        addBanpQueue    workqueue.TypedRateLimitingInterface[string]
308
        updateBanpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
309
        deleteBanpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.BaselineAdminNetworkPolicy]
310
        banpKeyMutex    keymutex.KeyMutex
311

312
        cnpsLister     anplisterv1alpha2.ClusterNetworkPolicyLister
313
        cnpsSynced     cache.InformerSynced
314
        addCnpQueue    workqueue.TypedRateLimitingInterface[string]
315
        updateCnpQueue workqueue.TypedRateLimitingInterface[*ClusterNetworkPolicyChangedDelta]
316
        deleteCnpQueue workqueue.TypedRateLimitingInterface[*netpolv1alpha2.ClusterNetworkPolicy]
317
        cnpKeyMutex    keymutex.KeyMutex
318

319
        csrLister           certListerv1.CertificateSigningRequestLister
320
        csrSynced           cache.InformerSynced
321
        addOrUpdateCsrQueue workqueue.TypedRateLimitingInterface[string]
322

323
        addOrUpdateVMIMigrationQueue workqueue.TypedRateLimitingInterface[string]
324
        deleteVMQueue                workqueue.TypedRateLimitingInterface[string]
325
        kubevirtInformerFactory      informer.KubeVirtInformerFactory
326

327
        netAttachLister          netAttachv1.NetworkAttachmentDefinitionLister
328
        netAttachSynced          cache.InformerSynced
329
        netAttachInformerFactory netAttach.SharedInformerFactory
330

331
        serviceCIDRStore           *util.ServiceCIDRStore
332
        serviceCIDRLister          netv1.ServiceCIDRLister
333
        serviceCIDRSynced          cache.InformerSynced
334
        serviceCIDRInformerFactory kubeinformers.SharedInformerFactory
335

336
        recorder               record.EventRecorder
337
        informerFactory        kubeinformers.SharedInformerFactory
338
        cmInformerFactory      kubeinformers.SharedInformerFactory
339
        deployInformerFactory  kubeinformers.SharedInformerFactory
340
        kubeovnInformerFactory kubeovninformer.SharedInformerFactory
341
        anpInformerFactory     anpinformer.SharedInformerFactory
342

343
        // Database health check
344
        dbFailureCount int
345

346
        distributedSubnetNeedSync atomic.Bool
347
}
348

349
func newTypedRateLimitingQueue[T comparable](name string, rateLimiter workqueue.TypedRateLimiter[T]) workqueue.TypedRateLimitingInterface[T] {
1✔
350
        if rateLimiter == nil {
2✔
351
                rateLimiter = workqueue.DefaultTypedControllerRateLimiter[T]()
1✔
352
        }
1✔
353
        return workqueue.NewTypedRateLimitingQueueWithConfig(rateLimiter, workqueue.TypedRateLimitingQueueConfig[T]{Name: name})
1✔
354
}
355

356
// Run creates and runs a new ovn controller
357
func Run(ctx context.Context, config *Configuration) {
×
358
        klog.V(4).Info("Creating event broadcaster")
×
359
        eventBroadcaster := record.NewBroadcasterWithCorrelatorOptions(record.CorrelatorOptions{BurstSize: 100})
×
360
        eventBroadcaster.StartLogging(klog.Infof)
×
361
        eventBroadcaster.StartRecordingToSink(&typedcorev1.EventSinkImpl{Interface: config.KubeFactoryClient.CoreV1().Events(metav1.NamespaceAll)})
×
362
        recorder := eventBroadcaster.NewRecorder(scheme.Scheme, corev1.EventSource{Component: controllerAgentName})
×
363
        custCrdRateLimiter := workqueue.NewTypedMaxOfRateLimiter(
×
364
                workqueue.NewTypedItemExponentialFailureRateLimiter[string](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
365
                &workqueue.TypedBucketRateLimiter[string]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
366
        )
×
367

×
368
        var err error
×
369
        informerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
370
                kubeinformers.WithTransform(util.TrimPodForController),
×
371
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
372
                        listOption.AllowWatchBookmarks = true
×
373
                }))
×
374
        cmInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
375
                kubeinformers.WithNamespace(config.PodNamespace),
×
376
                kubeinformers.WithTransform(util.TrimManagedFields),
×
377
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
378
                        listOption.AllowWatchBookmarks = true
×
379
                }))
×
380
        // deployment informer used to list/watch vpc egress gateway and nat gateway workloads
381
        deployInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
382
                kubeinformers.WithTransform(util.TrimManagedFields),
×
383
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
384
                        listOption.AllowWatchBookmarks = true
×
385
                }))
×
386
        kubeovnInformerFactory := kubeovninformer.NewSharedInformerFactoryWithOptions(config.KubeOvnFactoryClient, 0,
×
387
                kubeovninformer.WithTransform(util.TrimManagedFields),
×
388
                kubeovninformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
389
                        listOption.AllowWatchBookmarks = true
×
390
                }))
×
391
        anpInformerFactory := anpinformer.NewSharedInformerFactoryWithOptions(config.AnpClient, 0,
×
392
                anpinformer.WithTransform(util.TrimManagedFields),
×
393
                anpinformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
394
                        listOption.AllowWatchBookmarks = true
×
395
                }))
×
396
        attachNetInformerFactory := netAttach.NewSharedInformerFactoryWithOptions(config.AttachNetClient, 0,
×
397
                netAttach.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
398
                        listOption.AllowWatchBookmarks = true
×
399
                }),
×
400
        )
401
        kubevirtInformerFactory := informer.NewKubeVirtInformerFactoryWithOptions(config.KubevirtClient.RestClient(), config.KubevirtClient,
×
402
                informer.WithTransform(util.TrimManagedFields),
×
403
        )
×
404
        // Dedicated factory so that on clusters without the ServiceCIDR API the
×
405
        // failed list/watch does not contaminate the main informer factory.
×
406
        serviceCIDRInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeClient, 0,
×
407
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
408
                        listOption.AllowWatchBookmarks = true
×
409
                }),
×
410
        )
411

412
        vpcInformer := kubeovnInformerFactory.Kubeovn().V1().Vpcs()
×
413
        vpcNatGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcNatGateways()
×
414
        vpcEgressGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcEgressGateways()
×
415
        // BgpConf/EvpnConf informers are started lazily via StartBgpEvpnConfInformerFactory
×
416
        // because their CRDs are optional on clusters that don't use vpc-egress-gateway BGP/EVPN.
×
417
        subnetInformer := kubeovnInformerFactory.Kubeovn().V1().Subnets()
×
418
        ippoolInformer := kubeovnInformerFactory.Kubeovn().V1().IPPools()
×
419
        ipInformer := kubeovnInformerFactory.Kubeovn().V1().IPs()
×
420
        virtualIPInformer := kubeovnInformerFactory.Kubeovn().V1().Vips()
×
421
        iptablesEipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesEIPs()
×
422
        iptablesFipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesFIPRules()
×
423
        iptablesDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesDnatRules()
×
424
        iptablesSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesSnatRules()
×
425
        vlanInformer := kubeovnInformerFactory.Kubeovn().V1().Vlans()
×
426
        providerNetworkInformer := kubeovnInformerFactory.Kubeovn().V1().ProviderNetworks()
×
427
        sgInformer := kubeovnInformerFactory.Kubeovn().V1().SecurityGroups()
×
428
        podInformer := informerFactory.Core().V1().Pods()
×
429
        namespaceInformer := informerFactory.Core().V1().Namespaces()
×
430
        nodeInformer := informerFactory.Core().V1().Nodes()
×
431
        serviceInformer := informerFactory.Core().V1().Services()
×
432
        endpointSliceInformer := informerFactory.Discovery().V1().EndpointSlices()
×
433
        deploymentInformer := deployInformerFactory.Apps().V1().Deployments()
×
434
        qosPolicyInformer := kubeovnInformerFactory.Kubeovn().V1().QoSPolicies()
×
435
        configMapInformer := cmInformerFactory.Core().V1().ConfigMaps()
×
436
        npInformer := informerFactory.Networking().V1().NetworkPolicies()
×
437
        routerLBRuleInformer := kubeovnInformerFactory.Kubeovn().V1().RouterLBRules()
×
438
        switchLBRuleInformer := kubeovnInformerFactory.Kubeovn().V1().SwitchLBRules()
×
439
        vpcDNSInformer := kubeovnInformerFactory.Kubeovn().V1().VpcDnses()
×
440
        ovnEipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnEips()
×
441
        ovnFipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnFips()
×
442
        ovnSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnSnatRules()
×
443
        ovnDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnDnatRules()
×
444
        anpInformer := anpInformerFactory.Policy().V1alpha1().AdminNetworkPolicies()
×
445
        banpInformer := anpInformerFactory.Policy().V1alpha1().BaselineAdminNetworkPolicies()
×
446
        cnpInformer := anpInformerFactory.Policy().V1alpha2().ClusterNetworkPolicies()
×
447
        dnsNameResolverInformer := kubeovnInformerFactory.Kubeovn().V1().DNSNameResolvers()
×
448
        csrInformer := informerFactory.Certificates().V1().CertificateSigningRequests()
×
449
        netAttachInformer := attachNetInformerFactory.K8sCniCncfIo().V1().NetworkAttachmentDefinitions()
×
450

×
451
        numKeyLocks := max(runtime.NumCPU()*2, config.WorkerNum*2)
×
452
        controller := &Controller{
×
453
                config:             config,
×
454
                deletingPodObjMap:  xsync.NewMap[string, *corev1.Pod](),
×
455
                deletingNodeObjMap: xsync.NewMap[string, *corev1.Node](),
×
456
                ipam:               ovnipam.NewIPAM(),
×
457
                namedPort:          NewNamedPort(),
×
458

×
459
                vpcsLister:           vpcInformer.Lister(),
×
460
                vpcSynced:            vpcInformer.Informer().HasSynced,
×
461
                addOrUpdateVpcQueue:  newTypedRateLimitingQueue[string]("AddOrUpdateVpc", nil),
×
462
                vpcLastPoliciesMap:   xsync.NewMap[string, string](),
×
463
                delVpcQueue:          newTypedRateLimitingQueue[*kubeovnv1.Vpc]("DeleteVpc", nil),
×
464
                updateVpcStatusQueue: newTypedRateLimitingQueue[string]("UpdateVpcStatus", nil),
×
465
                vpcKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
466

×
467
                vpcNatGatewayLister:              vpcNatGatewayInformer.Lister(),
×
468
                vpcNatGatewaySynced:              vpcNatGatewayInformer.Informer().HasSynced,
×
469
                addOrUpdateVpcNatGatewayQueue:    newTypedRateLimitingQueue("AddOrUpdateVpcNatGw", custCrdRateLimiter),
×
470
                initVpcNatGatewayQueue:           newTypedRateLimitingQueue("InitVpcNatGw", custCrdRateLimiter),
×
471
                delVpcNatGatewayQueue:            newTypedRateLimitingQueue("DeleteVpcNatGw", custCrdRateLimiter),
×
472
                updateVpcEipQueue:                newTypedRateLimitingQueue("UpdateVpcEip", custCrdRateLimiter),
×
473
                updateVpcFloatingIPQueue:         newTypedRateLimitingQueue("UpdateVpcFloatingIp", custCrdRateLimiter),
×
474
                updateVpcDnatQueue:               newTypedRateLimitingQueue("UpdateVpcDnat", custCrdRateLimiter),
×
475
                updateVpcSnatQueue:               newTypedRateLimitingQueue("UpdateVpcSnat", custCrdRateLimiter),
×
476
                updateVpcSubnetQueue:             newTypedRateLimitingQueue("UpdateVpcSubnet", custCrdRateLimiter),
×
477
                vpcNatGwKeyMutex:                 keymutex.NewHashed(numKeyLocks),
×
478
                vpcNatGwExecKeyMutex:             keymutex.NewHashed(numKeyLocks),
×
479
                vpcEgressGatewayLister:           vpcEgressGatewayInformer.Lister(),
×
480
                vpcEgressGatewaySynced:           vpcEgressGatewayInformer.Informer().HasSynced,
×
481
                addOrUpdateVpcEgressGatewayQueue: newTypedRateLimitingQueue("AddOrUpdateVpcEgressGateway", custCrdRateLimiter),
×
482
                delVpcEgressGatewayQueue:         newTypedRateLimitingQueue("DeleteVpcEgressGateway", custCrdRateLimiter),
×
483
                vpcEgressGatewayKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
484

×
485
                // bgpConfLister/bgpConfSynced/evpnConfLister/evpnConfSynced are populated lazily
×
486
                // in startBgpEvpnConfInformer once the matching CRDs are detected.
×
487

×
488
                subnetsLister:           subnetInformer.Lister(),
×
489
                subnetSynced:            subnetInformer.Informer().HasSynced,
×
490
                addOrUpdateSubnetQueue:  newTypedRateLimitingQueue[string]("AddSubnet", nil),
×
491
                deleteSubnetQueue:       newTypedRateLimitingQueue[*kubeovnv1.Subnet]("DeleteSubnet", nil),
×
492
                updateSubnetStatusQueue: newTypedRateLimitingQueue[string]("UpdateSubnetStatus", nil),
×
493
                syncVirtualPortsQueue:   newTypedRateLimitingQueue[string]("SyncVirtualPort", nil),
×
494
                subnetKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
495

×
496
                ippoolLister:            ippoolInformer.Lister(),
×
497
                ippoolSynced:            ippoolInformer.Informer().HasSynced,
×
498
                addOrUpdateIPPoolQueue:  newTypedRateLimitingQueue[string]("AddIPPool", nil),
×
499
                updateIPPoolStatusQueue: newTypedRateLimitingQueue[string]("UpdateIPPoolStatus", nil),
×
500
                deleteIPPoolQueue:       newTypedRateLimitingQueue[*kubeovnv1.IPPool]("DeleteIPPool", nil),
×
501
                ippoolKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
502

×
503
                ipsLister:     ipInformer.Lister(),
×
504
                ipSynced:      ipInformer.Informer().HasSynced,
×
505
                addIPQueue:    newTypedRateLimitingQueue[string]("AddIP", nil),
×
506
                updateIPQueue: newTypedRateLimitingQueue[string]("UpdateIP", nil),
×
507
                delIPQueue:    newTypedRateLimitingQueue[*kubeovnv1.IP]("DeleteIP", nil),
×
508

×
509
                virtualIpsLister:          virtualIPInformer.Lister(),
×
510
                virtualIpsSynced:          virtualIPInformer.Informer().HasSynced,
×
511
                addVirtualIPQueue:         newTypedRateLimitingQueue[string]("AddVirtualIP", nil),
×
512
                updateVirtualIPQueue:      newTypedRateLimitingQueue[string]("UpdateVirtualIP", nil),
×
513
                updateVirtualParentsQueue: newTypedRateLimitingQueue[string]("UpdateVirtualParents", nil),
×
514
                delVirtualIPQueue:         newTypedRateLimitingQueue[*kubeovnv1.Vip]("DeleteVirtualIP", nil),
×
515

×
516
                iptablesEipsLister:     iptablesEipInformer.Lister(),
×
517
                iptablesEipSynced:      iptablesEipInformer.Informer().HasSynced,
×
518
                addIptablesEipQueue:    newTypedRateLimitingQueue("AddIptablesEip", custCrdRateLimiter),
×
519
                updateIptablesEipQueue: newTypedRateLimitingQueue("UpdateIptablesEip", custCrdRateLimiter),
×
520
                resetIptablesEipQueue:  newTypedRateLimitingQueue("ResetIptablesEip", custCrdRateLimiter),
×
521
                delIptablesEipQueue:    newTypedRateLimitingQueue[*kubeovnv1.IptablesEIP]("DeleteIptablesEip", nil),
×
522

×
523
                iptablesFipsLister:     iptablesFipInformer.Lister(),
×
524
                iptablesFipSynced:      iptablesFipInformer.Informer().HasSynced,
×
525
                addIptablesFipQueue:    newTypedRateLimitingQueue("AddIptablesFip", custCrdRateLimiter),
×
526
                updateIptablesFipQueue: newTypedRateLimitingQueue("UpdateIptablesFip", custCrdRateLimiter),
×
527
                delIptablesFipQueue:    newTypedRateLimitingQueue("DeleteIptablesFip", custCrdRateLimiter),
×
528

×
529
                iptablesDnatRulesLister:     iptablesDnatRuleInformer.Lister(),
×
530
                iptablesDnatRuleSynced:      iptablesDnatRuleInformer.Informer().HasSynced,
×
531
                addIptablesDnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesDnatRule", custCrdRateLimiter),
×
532
                updateIptablesDnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesDnatRule", custCrdRateLimiter),
×
533
                delIptablesDnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesDnatRule", custCrdRateLimiter),
×
534

×
535
                iptablesSnatRulesLister:     iptablesSnatRuleInformer.Lister(),
×
536
                iptablesSnatRuleSynced:      iptablesSnatRuleInformer.Informer().HasSynced,
×
537
                addIptablesSnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesSnatRule", custCrdRateLimiter),
×
538
                updateIptablesSnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesSnatRule", custCrdRateLimiter),
×
539
                delIptablesSnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesSnatRule", custCrdRateLimiter),
×
540

×
541
                vlansLister:     vlanInformer.Lister(),
×
542
                vlanSynced:      vlanInformer.Informer().HasSynced,
×
543
                addVlanQueue:    newTypedRateLimitingQueue[string]("AddVlan", nil),
×
544
                delVlanQueue:    newTypedRateLimitingQueue[string]("DeleteVlan", nil),
×
545
                updateVlanQueue: newTypedRateLimitingQueue[string]("UpdateVlan", nil),
×
546
                vlanKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
547

×
548
                providerNetworksLister: providerNetworkInformer.Lister(),
×
549
                providerNetworkSynced:  providerNetworkInformer.Informer().HasSynced,
×
550

×
551
                podsLister:          podInformer.Lister(),
×
552
                podsSynced:          podInformer.Informer().HasSynced,
×
553
                addOrUpdatePodQueue: newTypedRateLimitingQueue[string]("AddOrUpdatePod", nil),
×
554
                deletePodQueue: workqueue.NewTypedRateLimitingQueueWithConfig(
×
555
                        workqueue.DefaultTypedControllerRateLimiter[string](),
×
556
                        workqueue.TypedRateLimitingQueueConfig[string]{
×
557
                                Name:          "DeletePod",
×
558
                                DelayingQueue: workqueue.NewTypedDelayingQueue[string](),
×
559
                        },
×
560
                ),
×
561
                updatePodSecurityQueue: newTypedRateLimitingQueue[string]("UpdatePodSecurity", nil),
×
562
                podKeyMutex:            keymutex.NewHashed(numKeyLocks),
×
563

×
564
                namespacesLister:  namespaceInformer.Lister(),
×
565
                namespacesSynced:  namespaceInformer.Informer().HasSynced,
×
566
                addNamespaceQueue: newTypedRateLimitingQueue[string]("AddNamespace", nil),
×
567
                nsKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
568

×
569
                nodesLister:     nodeInformer.Lister(),
×
570
                nodesSynced:     nodeInformer.Informer().HasSynced,
×
571
                addNodeQueue:    newTypedRateLimitingQueue[string]("AddNode", nil),
×
572
                updateNodeQueue: newTypedRateLimitingQueue[string]("UpdateNode", nil),
×
573
                deleteNodeQueue: newTypedRateLimitingQueue[string]("DeleteNode", nil),
×
574
                nodeKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
575

×
576
                servicesLister:     serviceInformer.Lister(),
×
577
                serviceSynced:      serviceInformer.Informer().HasSynced,
×
578
                addServiceQueue:    newTypedRateLimitingQueue[string]("AddService", nil),
×
579
                deleteServiceQueue: newTypedRateLimitingQueue[*vpcService]("DeleteService", nil),
×
580
                updateServiceQueue: newTypedRateLimitingQueue[*updateSvcObject]("UpdateService", nil),
×
581
                svcKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
582

×
583
                endpointSlicesLister:          endpointSliceInformer.Lister(),
×
584
                endpointSlicesSynced:          endpointSliceInformer.Informer().HasSynced,
×
585
                addOrUpdateEndpointSliceQueue: newTypedRateLimitingQueue[string]("UpdateEndpointSlice", nil),
×
586
                epKeyMutex:                    keymutex.NewHashed(numKeyLocks),
×
587

×
588
                deploymentsLister: deploymentInformer.Lister(),
×
589
                deploymentsSynced: deploymentInformer.Informer().HasSynced,
×
590

×
591
                qosPoliciesLister:    qosPolicyInformer.Lister(),
×
592
                qosPolicySynced:      qosPolicyInformer.Informer().HasSynced,
×
593
                addQoSPolicyQueue:    newTypedRateLimitingQueue("AddQoSPolicy", custCrdRateLimiter),
×
594
                updateQoSPolicyQueue: newTypedRateLimitingQueue("UpdateQoSPolicy", custCrdRateLimiter),
×
595
                delQoSPolicyQueue:    newTypedRateLimitingQueue("DeleteQoSPolicy", custCrdRateLimiter),
×
596

×
597
                configMapsLister: configMapInformer.Lister(),
×
598
                configMapsSynced: configMapInformer.Informer().HasSynced,
×
599

×
600
                sgKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
601
                sgsLister:          sgInformer.Lister(),
×
602
                sgSynced:           sgInformer.Informer().HasSynced,
×
603
                addOrUpdateSgQueue: newTypedRateLimitingQueue[string]("UpdateSecurityGroup", nil),
×
604
                delSgQueue:         newTypedRateLimitingQueue[string]("DeleteSecurityGroup", nil),
×
605
                syncSgPortsQueue:   newTypedRateLimitingQueue[string]("SyncSecurityGroupPorts", nil),
×
606

×
607
                ovnEipsLister:     ovnEipInformer.Lister(),
×
608
                ovnEipSynced:      ovnEipInformer.Informer().HasSynced,
×
609
                addOvnEipQueue:    newTypedRateLimitingQueue("AddOvnEip", custCrdRateLimiter),
×
610
                updateOvnEipQueue: newTypedRateLimitingQueue("UpdateOvnEip", custCrdRateLimiter),
×
611
                resetOvnEipQueue:  newTypedRateLimitingQueue("ResetOvnEip", custCrdRateLimiter),
×
612
                delOvnEipQueue:    newTypedRateLimitingQueue[*kubeovnv1.OvnEip]("DeleteOvnEip", nil),
×
613

×
614
                ovnFipsLister:     ovnFipInformer.Lister(),
×
615
                ovnFipSynced:      ovnFipInformer.Informer().HasSynced,
×
616
                addOvnFipQueue:    newTypedRateLimitingQueue("AddOvnFip", custCrdRateLimiter),
×
617
                updateOvnFipQueue: newTypedRateLimitingQueue("UpdateOvnFip", custCrdRateLimiter),
×
618
                delOvnFipQueue:    newTypedRateLimitingQueue("DeleteOvnFip", custCrdRateLimiter),
×
619

×
620
                ovnSnatRulesLister:     ovnSnatRuleInformer.Lister(),
×
621
                ovnSnatRuleSynced:      ovnSnatRuleInformer.Informer().HasSynced,
×
622
                addOvnSnatRuleQueue:    newTypedRateLimitingQueue("AddOvnSnatRule", custCrdRateLimiter),
×
623
                updateOvnSnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnSnatRule", custCrdRateLimiter),
×
624
                delOvnSnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnSnatRule", custCrdRateLimiter),
×
625

×
626
                ovnDnatRulesLister:     ovnDnatRuleInformer.Lister(),
×
627
                ovnDnatRuleSynced:      ovnDnatRuleInformer.Informer().HasSynced,
×
628
                addOvnDnatRuleQueue:    newTypedRateLimitingQueue("AddOvnDnatRule", custCrdRateLimiter),
×
629
                updateOvnDnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnDnatRule", custCrdRateLimiter),
×
630
                delOvnDnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnDnatRule", custCrdRateLimiter),
×
631

×
632
                csrLister:           csrInformer.Lister(),
×
633
                csrSynced:           csrInformer.Informer().HasSynced,
×
634
                addOrUpdateCsrQueue: newTypedRateLimitingQueue("AddOrUpdateCSR", custCrdRateLimiter),
×
635

×
636
                addOrUpdateVMIMigrationQueue: newTypedRateLimitingQueue[string]("AddOrUpdateVMIMigration", nil),
×
637
                deleteVMQueue:                newTypedRateLimitingQueue[string]("DeleteVM", nil),
×
638
                kubevirtInformerFactory:      kubevirtInformerFactory,
×
639

×
640
                netAttachLister:          netAttachInformer.Lister(),
×
641
                netAttachSynced:          netAttachInformer.Informer().HasSynced,
×
642
                netAttachInformerFactory: attachNetInformerFactory,
×
643

×
644
                serviceCIDRStore:           util.NewServiceCIDRStore(config.ServiceClusterIPRange),
×
645
                serviceCIDRInformerFactory: serviceCIDRInformerFactory,
×
646

×
647
                recorder:               recorder,
×
648
                informerFactory:        informerFactory,
×
649
                cmInformerFactory:      cmInformerFactory,
×
650
                deployInformerFactory:  deployInformerFactory,
×
651
                kubeovnInformerFactory: kubeovnInformerFactory,
×
652
                anpInformerFactory:     anpInformerFactory,
×
653
        }
×
654

×
655
        if controller.OVNNbClient, err = ovs.NewOvnNbClient(
×
656
                config.OvnNbAddr,
×
657
                config.OvnTimeout,
×
658
                config.OvsDbConnectTimeout,
×
659
                config.OvsDbInactivityTimeout,
×
660
                config.OvsDbConnectMaxRetry,
×
661
        ); err != nil {
×
662
                util.LogFatalAndExit(err, "failed to create ovn nb client")
×
663
        }
×
664
        if controller.OVNSbClient, err = ovs.NewOvnSbClient(
×
665
                config.OvnSbAddr,
×
666
                config.OvnTimeout,
×
667
                config.OvsDbConnectTimeout,
×
668
                config.OvsDbInactivityTimeout,
×
669
                config.OvsDbConnectMaxRetry,
×
670
        ); err != nil {
×
671
                util.LogFatalAndExit(err, "failed to create ovn sb client")
×
672
        }
×
673
        if config.EnableLb {
×
674
                controller.routerLBRuleLister = routerLBRuleInformer.Lister()
×
675
                controller.routerLBRuleSynced = routerLBRuleInformer.Informer().HasSynced
×
676
                controller.addRouterLBRuleQueue = newTypedRateLimitingQueue("AddRouterLBRule", custCrdRateLimiter)
×
677
                controller.delRouterLBRuleQueue = newTypedRateLimitingQueue(
×
678
                        "DeleteRouterLBRule",
×
679
                        workqueue.NewTypedMaxOfRateLimiter(
×
680
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*RouterLBRuleInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
681
                                &workqueue.TypedBucketRateLimiter[*RouterLBRuleInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
682
                        ),
×
683
                )
×
684
                controller.updateRouterLBRuleQueue = newTypedRateLimitingQueue(
×
685
                        "UpdateRouterLBRule",
×
686
                        workqueue.NewTypedMaxOfRateLimiter(
×
687
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*RouterLBRuleInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
688
                                &workqueue.TypedBucketRateLimiter[*RouterLBRuleInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
689
                        ),
×
690
                )
×
691

×
692
                controller.switchLBRuleLister = switchLBRuleInformer.Lister()
×
693
                controller.switchLBRuleSynced = switchLBRuleInformer.Informer().HasSynced
×
694
                controller.addSwitchLBRuleQueue = newTypedRateLimitingQueue("AddSwitchLBRule", custCrdRateLimiter)
×
695
                controller.delSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
696
                        "DeleteSwitchLBRule",
×
697
                        workqueue.NewTypedMaxOfRateLimiter(
×
698
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SwitchLBRuleInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
699
                                &workqueue.TypedBucketRateLimiter[*SwitchLBRuleInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
700
                        ),
×
701
                )
×
702
                controller.updateSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
703
                        "UpdateSwitchLBRule",
×
704
                        workqueue.NewTypedMaxOfRateLimiter(
×
705
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SwitchLBRuleInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
706
                                &workqueue.TypedBucketRateLimiter[*SwitchLBRuleInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
707
                        ),
×
708
                )
×
709

×
710
                controller.vpcDNSLister = vpcDNSInformer.Lister()
×
711
                controller.vpcDNSSynced = vpcDNSInformer.Informer().HasSynced
×
712
                controller.addOrUpdateVpcDNSQueue = newTypedRateLimitingQueue("AddOrUpdateVpcDns", custCrdRateLimiter)
×
713
                controller.delVpcDNSQueue = newTypedRateLimitingQueue("DeleteVpcDns", custCrdRateLimiter)
×
714
        }
×
715

716
        if config.EnableNP {
×
717
                controller.npsLister = npInformer.Lister()
×
718
                controller.npsSynced = npInformer.Informer().HasSynced
×
719
                controller.npIndexer = npInformer.Informer().GetIndexer()
×
720
                controller.updateNpQueue = newTypedRateLimitingQueue[string]("UpdateNetworkPolicy", nil)
×
721
                controller.deleteNpQueue = newTypedRateLimitingQueue[string]("DeleteNetworkPolicy", nil)
×
722
                controller.npKeyMutex = keymutex.NewHashed(numKeyLocks)
×
723
        }
×
724

725
        if config.EnableANP {
×
726
                controller.anpsLister = anpInformer.Lister()
×
727
                controller.anpsSynced = anpInformer.Informer().HasSynced
×
728
                controller.addAnpQueue = newTypedRateLimitingQueue[string]("AddAdminNetworkPolicy", nil)
×
729
                controller.updateAnpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateAdminNetworkPolicy", nil)
×
730
                controller.deleteAnpQueue = newTypedRateLimitingQueue[*v1alpha1.AdminNetworkPolicy]("DeleteAdminNetworkPolicy", nil)
×
731
                controller.anpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
732

×
733
                controller.banpsLister = banpInformer.Lister()
×
734
                controller.banpsSynced = banpInformer.Informer().HasSynced
×
735
                controller.addBanpQueue = newTypedRateLimitingQueue[string]("AddBaseAdminNetworkPolicy", nil)
×
736
                controller.updateBanpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateBaseAdminNetworkPolicy", nil)
×
737
                controller.deleteBanpQueue = newTypedRateLimitingQueue[*v1alpha1.BaselineAdminNetworkPolicy]("DeleteBaseAdminNetworkPolicy", nil)
×
738
                controller.banpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
739

×
740
                controller.cnpsLister = cnpInformer.Lister()
×
741
                controller.cnpsSynced = cnpInformer.Informer().HasSynced
×
742
                controller.addCnpQueue = newTypedRateLimitingQueue[string]("AddClusterNetworkPolicy", nil)
×
743
                controller.updateCnpQueue = newTypedRateLimitingQueue[*ClusterNetworkPolicyChangedDelta]("UpdateClusterNetworkPolicy", nil)
×
744
                controller.deleteCnpQueue = newTypedRateLimitingQueue[*netpolv1alpha2.ClusterNetworkPolicy]("DeleteClusterNetworkPolicy", nil)
×
745
                controller.cnpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
746
        }
×
747

748
        if config.EnableDNSNameResolver {
×
749
                if !config.EnableANP {
×
750
                        klog.Warning("DNS name resolver is enabled but ANP support is disabled, DNSNameResolver resources will not take effect")
×
751
                }
×
752
                controller.dnsNameResolversLister = dnsNameResolverInformer.Lister()
×
753
                controller.dnsNameResolversSynced = dnsNameResolverInformer.Informer().HasSynced
×
754
                if err := dnsNameResolverInformer.Informer().AddIndexers(cache.Indexers{
×
755
                        IndexDNSNameResolverByName: indexDNSNameResolverByName,
×
756
                }); err != nil {
×
757
                        util.LogFatalAndExit(err, "failed to add DNSNameResolver indexer")
×
758
                }
×
759
                controller.dnsNameResolverIndexer = dnsNameResolverInformer.Informer().GetIndexer()
×
760
                controller.addOrUpdateDNSNameResolverQueue = newTypedRateLimitingQueue[string]("AddOrUpdateDNSNameResolver", nil)
×
761
                controller.deleteDNSNameResolverQueue = newTypedRateLimitingQueue[*kubeovnv1.DNSNameResolver]("DeleteDNSNameResolver", nil)
×
762
        }
763

764
        if err := controller.setupIndexers(vpcInformer.Informer(), podInformer.Informer(), endpointSliceInformer.Informer(), ipInformer.Informer()); err != nil {
×
765
                util.LogFatalAndExit(err, "failed to set up informer indexers")
×
766
        }
×
767

768
        defer controller.shutdown()
×
769
        klog.Info("Starting OVN controller")
×
770

×
771
        // Start and sync NAD informer first, as many resources depend on NAD cache
×
772
        // NAD CRD is optional, so we check if it exists before starting the informer
×
773
        controller.StartNetAttachInformerFactory(ctx)
×
774

×
775
        // ServiceCIDR (networking.k8s.io/v1) is GA in K8s 1.33; older clusters
×
776
        // don't have the API at all. Best-effort start with periodic retry.
×
777
        controller.StartServiceCIDRInformerFactory(ctx)
×
778

×
779
        // BgpConf/EvpnConf are optional CRDs (v1.16.0+, used by vpc-egress-gateway BGP/EVPN).
×
780
        // They may be missing on clusters upgraded from <v1.16 via Helm, which does not
×
781
        // re-apply the `crds/` directory on `helm upgrade`. Best-effort start with periodic retry.
×
782
        controller.StartBgpEvpnConfInformerFactory(ctx)
×
783

×
NEW
784
        // MetalLB is optional. When its ServiceL2Status API is available, use it to
×
NEW
785
        // identify the chassis announcing each underlay LoadBalancer VIP.
×
NEW
786
        controller.StartServiceL2StatusInformer(ctx)
×
NEW
787

×
788
        // Wait for the caches to be synced before starting workers
×
789
        controller.informerFactory.Start(ctx.Done())
×
790
        controller.cmInformerFactory.Start(ctx.Done())
×
791
        controller.deployInformerFactory.Start(ctx.Done())
×
792
        controller.kubeovnInformerFactory.Start(ctx.Done())
×
793
        controller.anpInformerFactory.Start(ctx.Done())
×
794
        controller.StartKubevirtInformerFactory(ctx, kubevirtInformerFactory)
×
795

×
796
        klog.Info("Waiting for informer caches to sync")
×
797
        cacheSyncs := []cache.InformerSynced{
×
798
                controller.vpcNatGatewaySynced, controller.vpcEgressGatewaySynced,
×
799
                controller.vpcSynced, controller.subnetSynced,
×
800
                controller.ipSynced, controller.virtualIpsSynced, controller.iptablesEipSynced,
×
801
                controller.iptablesFipSynced, controller.iptablesDnatRuleSynced, controller.iptablesSnatRuleSynced,
×
802
                controller.vlanSynced, controller.podsSynced, controller.namespacesSynced, controller.nodesSynced,
×
803
                controller.serviceSynced, controller.endpointSlicesSynced, controller.deploymentsSynced, controller.configMapsSynced,
×
804
                controller.ovnEipSynced, controller.ovnFipSynced, controller.ovnSnatRuleSynced,
×
805
                controller.ovnDnatRuleSynced,
×
806
        }
×
807
        if controller.config.EnableLb {
×
808
                cacheSyncs = append(cacheSyncs, controller.routerLBRuleSynced, controller.switchLBRuleSynced, controller.vpcDNSSynced)
×
809
        }
×
810
        if controller.config.EnableNP {
×
811
                cacheSyncs = append(cacheSyncs, controller.npsSynced)
×
812
        }
×
813
        if controller.config.EnableANP {
×
814
                cacheSyncs = append(cacheSyncs, controller.anpsSynced, controller.banpsSynced, controller.cnpsSynced)
×
815
        }
×
816
        if controller.config.EnableDNSNameResolver {
×
817
                cacheSyncs = append(cacheSyncs, controller.dnsNameResolversSynced)
×
818
        }
×
819

820
        if !cache.WaitForCacheSync(ctx.Done(), cacheSyncs...) {
×
821
                util.LogFatalAndExit(nil, "failed to wait for caches to sync")
×
822
        }
×
823

824
        if _, err = podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
825
                AddFunc:    controller.enqueueAddPod,
×
826
                DeleteFunc: controller.enqueueDeletePod,
×
827
                UpdateFunc: controller.enqueueUpdatePod,
×
828
        }); err != nil {
×
829
                util.LogFatalAndExit(err, "failed to add pod event handler")
×
830
        }
×
831

832
        if _, err = namespaceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
833
                AddFunc:    controller.enqueueAddNamespace,
×
834
                UpdateFunc: controller.enqueueUpdateNamespace,
×
835
                DeleteFunc: controller.enqueueDeleteNamespace,
×
836
        }); err != nil {
×
837
                util.LogFatalAndExit(err, "failed to add namespace event handler")
×
838
        }
×
839

840
        if _, err = nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
841
                AddFunc:    controller.enqueueAddNode,
×
842
                UpdateFunc: controller.enqueueUpdateNode,
×
843
                DeleteFunc: controller.enqueueDeleteNode,
×
844
        }); err != nil {
×
845
                util.LogFatalAndExit(err, "failed to add node event handler")
×
846
        }
×
847

848
        if _, err = serviceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
849
                AddFunc:    controller.enqueueAddService,
×
850
                DeleteFunc: controller.enqueueDeleteService,
×
851
                UpdateFunc: controller.enqueueUpdateService,
×
852
        }); err != nil {
×
853
                util.LogFatalAndExit(err, "failed to add service event handler")
×
854
        }
×
855

856
        if _, err = endpointSliceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
857
                AddFunc:    controller.enqueueAddEndpointSlice,
×
858
                UpdateFunc: controller.enqueueUpdateEndpointSlice,
×
859
        }); err != nil {
×
860
                util.LogFatalAndExit(err, "failed to add endpoint slice event handler")
×
861
        }
×
862

863
        if _, err = deploymentInformer.Informer().AddEventHandler(cache.FilteringResourceEventHandler{
×
864
                FilterFunc: func(obj any) bool {
×
865
                        if deploy, ok := obj.(*appsv1api.Deployment); ok {
×
866
                                // Only watch deployments with VpcEgressGatewayLabel or VpcNatGatewayLabel
×
867
                                _, hasNatGwLabel := deploy.Labels[util.VpcNatGatewayLabel]
×
868
                                _, hasEgressGwLabel := deploy.Labels[util.VpcEgressGatewayLabel]
×
869
                                return hasNatGwLabel || hasEgressGwLabel
×
870
                        }
×
871
                        return false
×
872
                },
873
                Handler: cache.ResourceEventHandlerFuncs{
874
                        AddFunc:    controller.enqueueAddDeployment,
875
                        UpdateFunc: controller.enqueueUpdateDeployment,
876
                },
877
        }); err != nil {
×
878
                util.LogFatalAndExit(err, "failed to add deployment event handler")
×
879
        }
×
880

881
        if _, err = vpcInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
882
                AddFunc:    controller.enqueueAddVpc,
×
883
                UpdateFunc: controller.enqueueUpdateVpc,
×
884
                DeleteFunc: controller.enqueueDelVpc,
×
885
        }); err != nil {
×
886
                util.LogFatalAndExit(err, "failed to add vpc event handler")
×
887
        }
×
888

889
        if _, err = vpcNatGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
890
                AddFunc:    controller.enqueueAddVpcNatGw,
×
891
                UpdateFunc: controller.enqueueUpdateVpcNatGw,
×
892
                DeleteFunc: controller.enqueueDeleteVpcNatGw,
×
893
        }); err != nil {
×
894
                util.LogFatalAndExit(err, "failed to add vpc nat gateway event handler")
×
895
        }
×
896

897
        if _, err = vpcEgressGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
898
                AddFunc:    controller.enqueueAddVpcEgressGateway,
×
899
                UpdateFunc: controller.enqueueUpdateVpcEgressGateway,
×
900
                DeleteFunc: controller.enqueueDeleteVpcEgressGateway,
×
901
        }); err != nil {
×
902
                util.LogFatalAndExit(err, "failed to add vpc egress gateway event handler")
×
903
        }
×
904

905
        if _, err = subnetInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
906
                AddFunc:    controller.enqueueAddSubnet,
×
907
                UpdateFunc: controller.enqueueUpdateSubnet,
×
908
                DeleteFunc: controller.enqueueDeleteSubnet,
×
909
        }); err != nil {
×
910
                util.LogFatalAndExit(err, "failed to add subnet event handler")
×
911
        }
×
912

913
        if _, err = ippoolInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
914
                AddFunc:    controller.enqueueAddIPPool,
×
915
                UpdateFunc: controller.enqueueUpdateIPPool,
×
916
                DeleteFunc: controller.enqueueDeleteIPPool,
×
917
        }); err != nil {
×
918
                util.LogFatalAndExit(err, "failed to add ippool event handler")
×
919
        }
×
920

921
        if _, err = ipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
922
                AddFunc:    controller.enqueueAddIP,
×
923
                UpdateFunc: controller.enqueueUpdateIP,
×
924
                DeleteFunc: controller.enqueueDelIP,
×
925
        }); err != nil {
×
926
                util.LogFatalAndExit(err, "failed to add ips event handler")
×
927
        }
×
928

929
        if _, err = vlanInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
930
                AddFunc:    controller.enqueueAddVlan,
×
931
                DeleteFunc: controller.enqueueDelVlan,
×
932
                UpdateFunc: controller.enqueueUpdateVlan,
×
933
        }); err != nil {
×
934
                util.LogFatalAndExit(err, "failed to add vlan event handler")
×
935
        }
×
936

937
        if _, err = sgInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
938
                AddFunc:    controller.enqueueAddSg,
×
939
                DeleteFunc: controller.enqueueDeleteSg,
×
940
                UpdateFunc: controller.enqueueUpdateSg,
×
941
        }); err != nil {
×
942
                util.LogFatalAndExit(err, "failed to add security group event handler")
×
943
        }
×
944

945
        if _, err = virtualIPInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
946
                AddFunc:    controller.enqueueAddVirtualIP,
×
947
                UpdateFunc: controller.enqueueUpdateVirtualIP,
×
948
                DeleteFunc: controller.enqueueDelVirtualIP,
×
949
        }); err != nil {
×
950
                util.LogFatalAndExit(err, "failed to add virtual ip event handler")
×
951
        }
×
952

953
        if _, err = iptablesEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
954
                AddFunc:    controller.enqueueAddIptablesEip,
×
955
                UpdateFunc: controller.enqueueUpdateIptablesEip,
×
956
                DeleteFunc: controller.enqueueDelIptablesEip,
×
957
        }); err != nil {
×
958
                util.LogFatalAndExit(err, "failed to add iptables eip event handler")
×
959
        }
×
960

961
        if _, err = iptablesFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
962
                AddFunc:    controller.enqueueAddIptablesFip,
×
963
                UpdateFunc: controller.enqueueUpdateIptablesFip,
×
964
                DeleteFunc: controller.enqueueDelIptablesFip,
×
965
        }); err != nil {
×
966
                util.LogFatalAndExit(err, "failed to add iptables fip event handler")
×
967
        }
×
968

969
        if _, err = iptablesDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
970
                AddFunc:    controller.enqueueAddIptablesDnatRule,
×
971
                UpdateFunc: controller.enqueueUpdateIptablesDnatRule,
×
972
                DeleteFunc: controller.enqueueDelIptablesDnatRule,
×
973
        }); err != nil {
×
974
                util.LogFatalAndExit(err, "failed to add iptables dnat event handler")
×
975
        }
×
976

977
        if _, err = iptablesSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
978
                AddFunc:    controller.enqueueAddIptablesSnatRule,
×
979
                UpdateFunc: controller.enqueueUpdateIptablesSnatRule,
×
980
                DeleteFunc: controller.enqueueDelIptablesSnatRule,
×
981
        }); err != nil {
×
982
                util.LogFatalAndExit(err, "failed to add iptables snat rule event handler")
×
983
        }
×
984

985
        if _, err = ovnEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
986
                AddFunc:    controller.enqueueAddOvnEip,
×
987
                UpdateFunc: controller.enqueueUpdateOvnEip,
×
988
                DeleteFunc: controller.enqueueDelOvnEip,
×
989
        }); err != nil {
×
990
                util.LogFatalAndExit(err, "failed to add ovn eip event handler")
×
991
        }
×
992

993
        if _, err = ovnFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
994
                AddFunc:    controller.enqueueAddOvnFip,
×
995
                UpdateFunc: controller.enqueueUpdateOvnFip,
×
996
                DeleteFunc: controller.enqueueDelOvnFip,
×
997
        }); err != nil {
×
998
                util.LogFatalAndExit(err, "failed to add ovn fip event handler")
×
999
        }
×
1000

1001
        if _, err = ovnSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1002
                AddFunc:    controller.enqueueAddOvnSnatRule,
×
1003
                UpdateFunc: controller.enqueueUpdateOvnSnatRule,
×
1004
                DeleteFunc: controller.enqueueDelOvnSnatRule,
×
1005
        }); err != nil {
×
1006
                util.LogFatalAndExit(err, "failed to add ovn snat rule event handler")
×
1007
        }
×
1008

1009
        if _, err = ovnDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1010
                AddFunc:    controller.enqueueAddOvnDnatRule,
×
1011
                UpdateFunc: controller.enqueueUpdateOvnDnatRule,
×
1012
                DeleteFunc: controller.enqueueDelOvnDnatRule,
×
1013
        }); err != nil {
×
1014
                util.LogFatalAndExit(err, "failed to add ovn dnat rule event handler")
×
1015
        }
×
1016

1017
        if _, err = qosPolicyInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1018
                AddFunc:    controller.enqueueAddQoSPolicy,
×
1019
                UpdateFunc: controller.enqueueUpdateQoSPolicy,
×
1020
                DeleteFunc: controller.enqueueDelQoSPolicy,
×
1021
        }); err != nil {
×
1022
                util.LogFatalAndExit(err, "failed to add qos policy event handler")
×
1023
        }
×
1024

1025
        if config.EnableLb {
×
1026
                if _, err = routerLBRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1027
                        AddFunc:    controller.enqueueAddRouterLBRule,
×
1028
                        UpdateFunc: controller.enqueueUpdateRouterLBRule,
×
1029
                        DeleteFunc: controller.enqueueDeleteRouterLBRule,
×
1030
                }); err != nil {
×
1031
                        util.LogFatalAndExit(err, "failed to add router lb rule event handler")
×
1032
                }
×
1033

1034
                if _, err = switchLBRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1035
                        AddFunc:    controller.enqueueAddSwitchLBRule,
×
1036
                        UpdateFunc: controller.enqueueUpdateSwitchLBRule,
×
1037
                        DeleteFunc: controller.enqueueDeleteSwitchLBRule,
×
1038
                }); err != nil {
×
1039
                        util.LogFatalAndExit(err, "failed to add switch lb rule event handler")
×
1040
                }
×
1041

1042
                if _, err = vpcDNSInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1043
                        AddFunc:    controller.enqueueAddVpcDNS,
×
1044
                        UpdateFunc: controller.enqueueUpdateVpcDNS,
×
1045
                        DeleteFunc: controller.enqueueDeleteVPCDNS,
×
1046
                }); err != nil {
×
1047
                        util.LogFatalAndExit(err, "failed to add vpc dns event handler")
×
1048
                }
×
1049
        }
1050

1051
        if config.EnableNP {
×
1052
                if _, err = npInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1053
                        AddFunc:    controller.enqueueAddNp,
×
1054
                        UpdateFunc: controller.enqueueUpdateNp,
×
1055
                        DeleteFunc: controller.enqueueDeleteNp,
×
1056
                }); err != nil {
×
1057
                        util.LogFatalAndExit(err, "failed to add network policy event handler")
×
1058
                }
×
1059
        }
1060

1061
        if config.EnableANP {
×
1062
                if _, err = anpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1063
                        AddFunc:    controller.enqueueAddAnp,
×
1064
                        UpdateFunc: controller.enqueueUpdateAnp,
×
1065
                        DeleteFunc: controller.enqueueDeleteAnp,
×
1066
                }); err != nil {
×
1067
                        util.LogFatalAndExit(err, "failed to add admin network policy event handler")
×
1068
                }
×
1069

1070
                if _, err = banpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1071
                        AddFunc:    controller.enqueueAddBanp,
×
1072
                        UpdateFunc: controller.enqueueUpdateBanp,
×
1073
                        DeleteFunc: controller.enqueueDeleteBanp,
×
1074
                }); err != nil {
×
1075
                        util.LogFatalAndExit(err, "failed to add baseline admin network policy event handler")
×
1076
                }
×
1077

1078
                if _, err = cnpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1079
                        AddFunc:    controller.enqueueAddCnp,
×
1080
                        UpdateFunc: controller.enqueueUpdateCnp,
×
1081
                        DeleteFunc: controller.enqueueDeleteCnp,
×
1082
                }); err != nil {
×
1083
                        util.LogFatalAndExit(err, "failed to add cluster network policy event handler")
×
1084
                }
×
1085

1086
                maxPriorityPerMap := util.CnpMaxPriority + 1
×
1087
                controller.anpPrioNameMap = make(map[int32]string, maxPriorityPerMap)
×
1088
                controller.anpNamePrioMap = make(map[string]int32, maxPriorityPerMap)
×
1089
                controller.bnpPrioNameMap = make(map[int32]string, maxPriorityPerMap)
×
1090
                controller.bnpNamePrioMap = make(map[string]int32, maxPriorityPerMap)
×
1091
        }
1092

1093
        if config.EnableDNSNameResolver {
×
1094
                if _, err = dnsNameResolverInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1095
                        AddFunc:    controller.enqueueAddDNSNameResolver,
×
1096
                        UpdateFunc: controller.enqueueUpdateDNSNameResolver,
×
1097
                        DeleteFunc: controller.enqueueDeleteDNSNameResolver,
×
1098
                }); err != nil {
×
1099
                        util.LogFatalAndExit(err, "failed to add dns name resolver event handler")
×
1100
                }
×
1101
        }
1102

1103
        if config.EnableOVNIPSec {
×
1104
                if _, err = csrInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1105
                        AddFunc:    controller.enqueueAddCsr,
×
1106
                        UpdateFunc: controller.enqueueUpdateCsr,
×
1107
                        // no need to add delete func for csr
×
1108
                }); err != nil {
×
1109
                        util.LogFatalAndExit(err, "failed to add csr event handler")
×
1110
                }
×
1111
        }
1112

1113
        controller.Run(ctx)
×
1114
}
1115

1116
// Run will set up the event handlers for types we are interested in, as well
1117
// as syncing informer caches and starting workers. It will block until stopCh
1118
// is closed, at which point it will shutdown the workqueue and wait for
1119
// workers to finish processing their current work items.
1120
func (c *Controller) Run(ctx context.Context) {
×
1121
        // 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...
×
1122
        // Otherwise, the init process should be placed after all workers have already started working
×
1123
        if err := c.OVNNbClient.SetLsDnatModDlDst(c.config.LsDnatModDlDst); err != nil {
×
1124
                util.LogFatalAndExit(err, "failed to set NB_Global option ls_dnat_mod_dl_dst")
×
1125
        }
×
1126

1127
        if err := c.OVNNbClient.SetUseCtInvMatch(); err != nil {
×
1128
                util.LogFatalAndExit(err, "failed to set NB_Global option use_ct_inv_match to false")
×
1129
        }
×
1130

1131
        if err := c.OVNNbClient.SetLsCtSkipDstLportIPs(c.config.LsCtSkipDstLportIPs); err != nil {
×
1132
                util.LogFatalAndExit(err, "failed to set NB_Global option ls_ct_skip_dst_lport_ips")
×
1133
        }
×
1134

1135
        if err := c.OVNNbClient.SetNodeLocalDNSIP(strings.Join(c.config.NodeLocalDNSIPs, ",")); err != nil {
×
1136
                util.LogFatalAndExit(err, "failed to set NB_Global option node_local_dns_ip")
×
1137
        }
×
1138

1139
        if err := c.OVNNbClient.SetSkipConntrackCidrs(c.config.SkipConntrackDstCidrs); err != nil {
×
1140
                util.LogFatalAndExit(err, "failed to set NB_Global option skip_conntrack_ipcidrs")
×
1141
        }
×
1142

1143
        if err := c.OVNNbClient.SetOVNIPSec(c.config.EnableOVNIPSec); err != nil {
×
1144
                util.LogFatalAndExit(err, "failed to set NB_Global ipsec")
×
1145
        }
×
1146

1147
        if err := c.InitOVN(); err != nil {
×
1148
                util.LogFatalAndExit(err, "failed to initialize ovn resources")
×
1149
        }
×
1150

1151
        // sync ip crd before initIPAM since ip crd will be used to restore vm and statefulset pod in initIPAM
1152
        if err := c.syncIPCR(); err != nil {
×
1153
                util.LogFatalAndExit(err, "failed to sync crd ips")
×
1154
        }
×
1155

1156
        if err := c.syncFinalizers(); err != nil {
×
1157
                util.LogFatalAndExit(err, "failed to initialize crd finalizers")
×
1158
        }
×
1159

1160
        if err := c.InitIPAM(); err != nil {
×
1161
                util.LogFatalAndExit(err, "failed to initialize ipam")
×
1162
        }
×
1163

1164
        if err := c.syncNodeRoutes(); err != nil {
×
1165
                util.LogFatalAndExit(err, "failed to initialize node routes")
×
1166
        }
×
1167

1168
        if err := c.syncSubnetCR(); err != nil {
×
1169
                util.LogFatalAndExit(err, "failed to sync crd subnets")
×
1170
        }
×
1171

1172
        if err := c.syncVlanCR(); err != nil {
×
1173
                util.LogFatalAndExit(err, "failed to sync crd vlans")
×
1174
        }
×
1175

1176
        if c.config.EnableOVNIPSec && !c.config.CertManagerIPSecCert {
×
1177
                if err := c.InitDefaultOVNIPsecCA(); err != nil {
×
1178
                        util.LogFatalAndExit(err, "failed to init ovn ipsec CA")
×
1179
                }
×
1180
        }
1181

1182
        c.startKubeOVNTLSManager(ctx)
×
1183

×
1184
        // start workers to do all the network operations
×
1185
        c.startWorkers(ctx)
×
1186

×
1187
        c.initResourceOnce()
×
1188
        <-ctx.Done()
×
1189
        klog.Info("Shutting down workers")
×
1190

×
1191
        c.OVNNbClient.Close()
×
1192
        c.OVNSbClient.Close()
×
1193
}
1194

1195
func (c *Controller) dbStatus() {
×
1196
        const maxFailures = 5
×
1197

×
1198
        done := make(chan error, 2)
×
1199
        go func() {
×
1200
                done <- c.OVNNbClient.Echo(context.Background())
×
1201
        }()
×
1202
        go func() {
×
1203
                done <- c.OVNSbClient.Echo(context.Background())
×
1204
        }()
×
1205

1206
        resultsReceived := 0
×
1207
        timeout := time.After(time.Duration(c.config.OvnTimeout) * time.Second)
×
1208

×
1209
        for resultsReceived < 2 {
×
1210
                select {
×
1211
                case err := <-done:
×
1212
                        resultsReceived++
×
1213
                        if err != nil {
×
1214
                                c.dbFailureCount++
×
1215
                                klog.Errorf("OVN database echo failed (%d/%d): %v", c.dbFailureCount, maxFailures, err)
×
1216
                                if c.dbFailureCount >= maxFailures {
×
1217
                                        util.LogFatalAndExit(err, "OVN database connection failed after %d attempts", maxFailures)
×
1218
                                }
×
1219
                                return
×
1220
                        }
1221
                case <-timeout:
×
1222
                        c.dbFailureCount++
×
1223
                        klog.Errorf("OVN database echo timeout (%d/%d) after %ds", c.dbFailureCount, maxFailures, c.config.OvnTimeout)
×
1224
                        if c.dbFailureCount >= maxFailures {
×
1225
                                util.LogFatalAndExit(nil, "OVN database connection timeout after %d attempts", maxFailures)
×
1226
                        }
×
1227
                        return
×
1228
                }
1229
        }
1230

1231
        if c.dbFailureCount > 0 {
×
1232
                klog.Infof("OVN database connection recovered after %d failures", c.dbFailureCount)
×
1233
                c.dbFailureCount = 0
×
1234
        }
×
1235
}
1236

1237
func (c *Controller) shutdown() {
×
1238
        utilruntime.HandleCrash()
×
1239

×
1240
        c.addOrUpdatePodQueue.ShutDown()
×
1241
        c.deletePodQueue.ShutDown()
×
1242
        c.updatePodSecurityQueue.ShutDown()
×
1243

×
1244
        c.addNamespaceQueue.ShutDown()
×
1245

×
1246
        c.addOrUpdateSubnetQueue.ShutDown()
×
1247
        c.deleteSubnetQueue.ShutDown()
×
1248
        c.updateSubnetStatusQueue.ShutDown()
×
1249
        c.syncVirtualPortsQueue.ShutDown()
×
1250

×
1251
        c.addOrUpdateIPPoolQueue.ShutDown()
×
1252
        c.updateIPPoolStatusQueue.ShutDown()
×
1253
        c.deleteIPPoolQueue.ShutDown()
×
1254

×
1255
        c.addNodeQueue.ShutDown()
×
1256
        c.updateNodeQueue.ShutDown()
×
1257
        c.deleteNodeQueue.ShutDown()
×
1258

×
1259
        c.addServiceQueue.ShutDown()
×
1260
        c.deleteServiceQueue.ShutDown()
×
1261
        c.updateServiceQueue.ShutDown()
×
1262
        c.addOrUpdateEndpointSliceQueue.ShutDown()
×
1263

×
1264
        c.addVlanQueue.ShutDown()
×
1265
        c.delVlanQueue.ShutDown()
×
1266
        c.updateVlanQueue.ShutDown()
×
1267

×
1268
        c.addOrUpdateVpcQueue.ShutDown()
×
1269
        c.updateVpcStatusQueue.ShutDown()
×
1270
        c.delVpcQueue.ShutDown()
×
1271

×
1272
        c.addOrUpdateVpcNatGatewayQueue.ShutDown()
×
1273
        c.initVpcNatGatewayQueue.ShutDown()
×
1274
        c.delVpcNatGatewayQueue.ShutDown()
×
1275
        c.updateVpcEipQueue.ShutDown()
×
1276
        c.updateVpcFloatingIPQueue.ShutDown()
×
1277
        c.updateVpcDnatQueue.ShutDown()
×
1278
        c.updateVpcSnatQueue.ShutDown()
×
1279
        c.updateVpcSubnetQueue.ShutDown()
×
1280

×
1281
        c.addOrUpdateVpcEgressGatewayQueue.ShutDown()
×
1282
        c.delVpcEgressGatewayQueue.ShutDown()
×
1283

×
1284
        if c.config.EnableLb {
×
1285
                c.addRouterLBRuleQueue.ShutDown()
×
1286
                c.delRouterLBRuleQueue.ShutDown()
×
1287
                c.updateRouterLBRuleQueue.ShutDown()
×
1288

×
1289
                c.addSwitchLBRuleQueue.ShutDown()
×
1290
                c.delSwitchLBRuleQueue.ShutDown()
×
1291
                c.updateSwitchLBRuleQueue.ShutDown()
×
1292

×
1293
                c.addOrUpdateVpcDNSQueue.ShutDown()
×
1294
                c.delVpcDNSQueue.ShutDown()
×
1295
        }
×
1296

1297
        c.addIPQueue.ShutDown()
×
1298
        c.updateIPQueue.ShutDown()
×
1299
        c.delIPQueue.ShutDown()
×
1300

×
1301
        c.addVirtualIPQueue.ShutDown()
×
1302
        c.updateVirtualIPQueue.ShutDown()
×
1303
        c.updateVirtualParentsQueue.ShutDown()
×
1304
        c.delVirtualIPQueue.ShutDown()
×
1305

×
1306
        c.addIptablesEipQueue.ShutDown()
×
1307
        c.updateIptablesEipQueue.ShutDown()
×
1308
        c.resetIptablesEipQueue.ShutDown()
×
1309
        c.delIptablesEipQueue.ShutDown()
×
1310

×
1311
        c.addIptablesFipQueue.ShutDown()
×
1312
        c.updateIptablesFipQueue.ShutDown()
×
1313
        c.delIptablesFipQueue.ShutDown()
×
1314

×
1315
        c.addIptablesDnatRuleQueue.ShutDown()
×
1316
        c.updateIptablesDnatRuleQueue.ShutDown()
×
1317
        c.delIptablesDnatRuleQueue.ShutDown()
×
1318

×
1319
        c.addIptablesSnatRuleQueue.ShutDown()
×
1320
        c.updateIptablesSnatRuleQueue.ShutDown()
×
1321
        c.delIptablesSnatRuleQueue.ShutDown()
×
1322

×
1323
        c.addQoSPolicyQueue.ShutDown()
×
1324
        c.updateQoSPolicyQueue.ShutDown()
×
1325
        c.delQoSPolicyQueue.ShutDown()
×
1326

×
1327
        c.addOvnEipQueue.ShutDown()
×
1328
        c.updateOvnEipQueue.ShutDown()
×
1329
        c.resetOvnEipQueue.ShutDown()
×
1330
        c.delOvnEipQueue.ShutDown()
×
1331

×
1332
        c.addOvnFipQueue.ShutDown()
×
1333
        c.updateOvnFipQueue.ShutDown()
×
1334
        c.delOvnFipQueue.ShutDown()
×
1335

×
1336
        c.addOvnSnatRuleQueue.ShutDown()
×
1337
        c.updateOvnSnatRuleQueue.ShutDown()
×
1338
        c.delOvnSnatRuleQueue.ShutDown()
×
1339

×
1340
        c.addOvnDnatRuleQueue.ShutDown()
×
1341
        c.updateOvnDnatRuleQueue.ShutDown()
×
1342
        c.delOvnDnatRuleQueue.ShutDown()
×
1343

×
1344
        if c.config.EnableNP {
×
1345
                c.updateNpQueue.ShutDown()
×
1346
                c.deleteNpQueue.ShutDown()
×
1347
        }
×
1348
        if c.config.EnableANP {
×
1349
                c.addAnpQueue.ShutDown()
×
1350
                c.updateAnpQueue.ShutDown()
×
1351
                c.deleteAnpQueue.ShutDown()
×
1352

×
1353
                c.addBanpQueue.ShutDown()
×
1354
                c.updateBanpQueue.ShutDown()
×
1355
                c.deleteBanpQueue.ShutDown()
×
1356

×
1357
                c.addCnpQueue.ShutDown()
×
1358
                c.updateCnpQueue.ShutDown()
×
1359
                c.deleteCnpQueue.ShutDown()
×
1360
        }
×
1361

1362
        if c.config.EnableDNSNameResolver {
×
1363
                c.addOrUpdateDNSNameResolverQueue.ShutDown()
×
1364
                c.deleteDNSNameResolverQueue.ShutDown()
×
1365
        }
×
1366

1367
        c.addOrUpdateSgQueue.ShutDown()
×
1368
        c.delSgQueue.ShutDown()
×
1369
        c.syncSgPortsQueue.ShutDown()
×
1370

×
1371
        c.addOrUpdateCsrQueue.ShutDown()
×
1372

×
1373
        if c.config.EnableLiveMigrationOptimize {
×
1374
                c.addOrUpdateVMIMigrationQueue.ShutDown()
×
1375
        }
×
1376
}
1377

1378
func (c *Controller) startWorkers(ctx context.Context) {
×
1379
        klog.Info("Starting workers")
×
1380

×
1381
        go wait.Until(runWorker("add/update vpc", c.addOrUpdateVpcQueue, c.handleAddOrUpdateVpc), time.Second, ctx.Done())
×
1382
        go wait.Until(runWorker("delete vpc", c.delVpcQueue, c.handleDelVpc), time.Second, ctx.Done())
×
1383
        go wait.Until(runWorker("update status of vpc", c.updateVpcStatusQueue, c.handleUpdateVpcStatus), time.Second, ctx.Done())
×
1384

×
1385
        go wait.Until(runWorker("add/update vpc nat gateway", c.addOrUpdateVpcNatGatewayQueue, c.handleAddOrUpdateVpcNatGw), time.Second, ctx.Done())
×
1386
        go wait.Until(runWorker("init vpc nat gateway", c.initVpcNatGatewayQueue, c.handleInitVpcNatGw), time.Second, ctx.Done())
×
1387
        go wait.Until(runWorker("delete vpc nat gateway", c.delVpcNatGatewayQueue, c.handleDelVpcNatGw), time.Second, ctx.Done())
×
1388
        go wait.Until(runWorker("add/update vpc egress gateway", c.addOrUpdateVpcEgressGatewayQueue, c.handleAddOrUpdateVpcEgressGateway), time.Second, ctx.Done())
×
1389
        go wait.Until(runWorker("delete vpc egress gateway", c.delVpcEgressGatewayQueue, c.handleDelVpcEgressGateway), time.Second, ctx.Done())
×
1390
        go wait.Until(runWorker("update fip for vpc nat gateway", c.updateVpcFloatingIPQueue, c.handleUpdateVpcFloatingIP), time.Second, ctx.Done())
×
1391
        go wait.Until(runWorker("update eip for vpc nat gateway", c.updateVpcEipQueue, c.handleUpdateVpcEip), time.Second, ctx.Done())
×
1392
        go wait.Until(runWorker("update dnat for vpc nat gateway", c.updateVpcDnatQueue, c.handleUpdateVpcDnat), time.Second, ctx.Done())
×
1393
        go wait.Until(runWorker("update snat for vpc nat gateway", c.updateVpcSnatQueue, c.handleUpdateVpcSnat), time.Second, ctx.Done())
×
1394
        go wait.Until(runWorker("update subnet route for vpc nat gateway", c.updateVpcSubnetQueue, c.handleUpdateNatGwSubnetRoute), time.Second, ctx.Done())
×
1395
        go wait.Until(runWorker("add/update csr", c.addOrUpdateCsrQueue, c.handleAddOrUpdateCsr), time.Second, ctx.Done())
×
1396
        // add default and join subnet and wait them ready
×
1397
        for range c.config.WorkerNum {
×
1398
                go wait.Until(runWorker("add/update subnet", c.addOrUpdateSubnetQueue, c.handleAddOrUpdateSubnet), time.Second, ctx.Done())
×
1399
        }
×
1400
        go wait.Until(runWorker("add/update ippool", c.addOrUpdateIPPoolQueue, c.handleAddOrUpdateIPPool), time.Second, ctx.Done())
×
1401
        go wait.Until(runWorker("add vlan", c.addVlanQueue, c.handleAddVlan), time.Second, ctx.Done())
×
1402
        go wait.Until(runWorker("add namespace", c.addNamespaceQueue, c.handleAddNamespace), time.Second, ctx.Done())
×
1403
        err := wait.PollUntilContextCancel(ctx, 3*time.Second, true, func(_ context.Context) (done bool, err error) {
×
1404
                subnets := []string{c.config.DefaultLogicalSwitch, c.config.NodeSwitch}
×
1405
                klog.Infof("wait for subnets %v ready", subnets)
×
1406

×
1407
                return c.allSubnetReady(subnets...)
×
1408
        })
×
1409
        if err != nil {
×
1410
                klog.Fatalf("wait default and join subnet ready, error: %v", err)
×
1411
        }
×
1412

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

×
1417
        // run node worker before handle any pods
×
1418
        for range c.config.WorkerNum {
×
1419
                go wait.Until(runWorker("add node", c.addNodeQueue, c.handleAddNode), time.Second, ctx.Done())
×
1420
                go wait.Until(runWorker("update node", c.updateNodeQueue, c.handleUpdateNode), time.Second, ctx.Done())
×
1421
                go wait.Until(runWorker("delete node", c.deleteNodeQueue, c.handleDeleteNode), time.Second, ctx.Done())
×
1422
        }
×
1423
        for {
×
1424
                ready := true
×
1425
                time.Sleep(3 * time.Second)
×
1426
                nodes, err := c.nodesLister.List(labels.Everything())
×
1427
                if err != nil {
×
1428
                        util.LogFatalAndExit(err, "failed to list nodes")
×
1429
                }
×
1430
                for _, node := range nodes {
×
1431
                        if node.Annotations[util.AllocatedAnnotation] != "true" {
×
1432
                                klog.Infof("wait node %s annotation ready", node.Name)
×
1433
                                ready = false
×
1434
                                break
×
1435
                        }
1436
                }
1437
                if ready {
×
1438
                        break
×
1439
                }
1440
        }
1441

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

×
1447
                go wait.Until(runWorker("add/update router lb rule", c.addRouterLBRuleQueue, c.handleAddOrUpdateRouterLBRule), time.Second, ctx.Done())
×
1448
                go wait.Until(runWorker("delete router lb rule", c.delRouterLBRuleQueue, c.handleDelRouterLBRule), time.Second, ctx.Done())
×
1449
                go wait.Until(runWorker("update router lb rule", c.updateRouterLBRuleQueue, c.handleUpdateRouterLBRule), time.Second, ctx.Done())
×
1450

×
1451
                go wait.Until(runWorker("add/update switch lb rule", c.addSwitchLBRuleQueue, c.handleAddOrUpdateSwitchLBRule), time.Second, ctx.Done())
×
1452
                go wait.Until(runWorker("delete switch lb rule", c.delSwitchLBRuleQueue, c.handleDelSwitchLBRule), time.Second, ctx.Done())
×
1453
                go wait.Until(runWorker("delete switch lb rule", c.updateSwitchLBRuleQueue, c.handleUpdateSwitchLBRule), time.Second, ctx.Done())
×
1454

×
1455
                go wait.Until(runWorker("add/update vpc dns", c.addOrUpdateVpcDNSQueue, c.handleAddOrUpdateVPCDNS), time.Second, ctx.Done())
×
1456
                go wait.Until(runWorker("delete vpc dns", c.delVpcDNSQueue, c.handleDelVpcDNS), time.Second, ctx.Done())
×
1457
                go wait.Until(func() {
×
1458
                        c.resyncVpcDNSConfig()
×
1459
                }, 5*time.Second, ctx.Done())
×
1460
        }
1461

1462
        for range c.config.WorkerNum {
×
1463
                go wait.Until(runWorker("delete pod", c.deletePodQueue, c.handleDeletePod), time.Second, ctx.Done())
×
1464
                go wait.Until(runWorker("add/update pod", c.addOrUpdatePodQueue, c.handleAddOrUpdatePod), time.Second, ctx.Done())
×
1465
                go wait.Until(runWorker("update pod security", c.updatePodSecurityQueue, c.handleUpdatePodSecurity), time.Second, ctx.Done())
×
1466

×
1467
                go wait.Until(runWorker("delete subnet", c.deleteSubnetQueue, c.handleDeleteSubnet), time.Second, ctx.Done())
×
1468
                go wait.Until(runWorker("delete ippool", c.deleteIPPoolQueue, c.handleDeleteIPPool), time.Second, ctx.Done())
×
1469
                go wait.Until(runWorker("update status of subnet", c.updateSubnetStatusQueue, c.handleUpdateSubnetStatus), time.Second, ctx.Done())
×
1470
                go wait.Until(runWorker("update status of ippool", c.updateIPPoolStatusQueue, c.handleUpdateIPPoolStatus), time.Second, ctx.Done())
×
1471
                go wait.Until(runWorker("virtual port for subnet", c.syncVirtualPortsQueue, c.syncVirtualPort), time.Second, ctx.Done())
×
1472

×
1473
                if c.config.EnableLb {
×
1474
                        go wait.Until(runWorker("update service", c.updateServiceQueue, c.handleUpdateService), time.Second, ctx.Done())
×
1475
                        go wait.Until(runWorker("add/update endpoint slice", c.addOrUpdateEndpointSliceQueue, c.handleUpdateEndpointSlice), time.Second, ctx.Done())
×
1476
                }
×
1477

1478
                if c.config.EnableNP {
×
1479
                        go wait.Until(runWorker("update network policy", c.updateNpQueue, c.handleUpdateNp), time.Second, ctx.Done())
×
1480
                        go wait.Until(runWorker("delete network policy", c.deleteNpQueue, c.handleDeleteNp), time.Second, ctx.Done())
×
1481
                }
×
1482

1483
                go wait.Until(runWorker("delete vlan", c.delVlanQueue, c.handleDelVlan), time.Second, ctx.Done())
×
1484
                go wait.Until(runWorker("update vlan", c.updateVlanQueue, c.handleUpdateVlan), time.Second, ctx.Done())
×
1485
        }
1486

1487
        if c.config.EnableEipSnat {
×
1488
                go wait.Until(func() {
×
1489
                        // init l3 about the default vpc external lrp binding to the gw chassis
×
1490
                        c.resyncExternalGateway()
×
1491
                }, time.Second, ctx.Done())
×
1492

1493
                // maintain l3 ha about the vpc external lrp binding to the gw chassis
1494
                c.OVNNbClient.MonitorBFD()
×
1495
        }
1496
        // TODO: we should merge these two vpc nat config into one config and resync them together
1497
        go wait.Until(func() {
×
1498
                c.resyncVpcNatGwConfig()
×
1499
        }, time.Second, ctx.Done())
×
1500

1501
        go wait.Until(func() {
×
1502
                c.resyncVpcNatConfig()
×
1503
        }, time.Second, ctx.Done())
×
1504

1505
        if c.config.GCInterval != 0 {
×
1506
                go wait.Until(func() {
×
1507
                        if err := c.markAndCleanLSP(); err != nil {
×
1508
                                klog.Errorf("gc lsp error: %v", err)
×
1509
                        }
×
1510
                }, time.Duration(c.config.GCInterval)*time.Second, ctx.Done())
1511
        }
1512

1513
        go wait.Until(func() {
×
1514
                if err := c.inspectPod(); err != nil {
×
1515
                        klog.Errorf("inspection error: %v", err)
×
1516
                }
×
1517
        }, time.Duration(c.config.InspectInterval)*time.Second, ctx.Done())
1518

1519
        if c.config.EnableExternalVpc {
×
1520
                go wait.Until(func() {
×
1521
                        c.syncExternalVpc()
×
1522
                }, 5*time.Second, ctx.Done())
×
1523
        }
1524

1525
        go wait.Until(c.resyncProviderNetworkStatus, 30*time.Second, ctx.Done())
×
1526
        go wait.Until(c.exportSubnetMetrics, 30*time.Second, ctx.Done())
×
1527
        go wait.Until(c.checkSubnetGateway, 5*time.Second, ctx.Done())
×
1528
        go wait.Until(c.syncDistributedSubnetRoutes, 5*time.Second, ctx.Done())
×
1529

×
1530
        go wait.Until(runWorker("add ovn eip", c.addOvnEipQueue, c.handleAddOvnEip), time.Second, ctx.Done())
×
1531
        go wait.Until(runWorker("update ovn eip", c.updateOvnEipQueue, c.handleUpdateOvnEip), time.Second, ctx.Done())
×
1532
        go wait.Until(runWorker("reset ovn eip", c.resetOvnEipQueue, c.handleResetOvnEip), time.Second, ctx.Done())
×
1533
        go wait.Until(runWorker("delete ovn eip", c.delOvnEipQueue, c.handleDelOvnEip), time.Second, ctx.Done())
×
1534

×
1535
        go wait.Until(runWorker("add ovn fip", c.addOvnFipQueue, c.handleAddOvnFip), time.Second, ctx.Done())
×
1536
        go wait.Until(runWorker("update ovn fip", c.updateOvnFipQueue, c.handleUpdateOvnFip), time.Second, ctx.Done())
×
1537
        go wait.Until(runWorker("delete ovn fip", c.delOvnFipQueue, c.handleDelOvnFip), time.Second, ctx.Done())
×
1538

×
1539
        go wait.Until(runWorker("add ovn snat rule", c.addOvnSnatRuleQueue, c.handleAddOvnSnatRule), time.Second, ctx.Done())
×
1540
        go wait.Until(runWorker("update ovn snat rule", c.updateOvnSnatRuleQueue, c.handleUpdateOvnSnatRule), time.Second, ctx.Done())
×
1541
        go wait.Until(runWorker("delete ovn snat rule", c.delOvnSnatRuleQueue, c.handleDelOvnSnatRule), time.Second, ctx.Done())
×
1542

×
1543
        go wait.Until(runWorker("add ovn dnat", c.addOvnDnatRuleQueue, c.handleAddOvnDnatRule), time.Second, ctx.Done())
×
1544
        go wait.Until(runWorker("update ovn dnat", c.updateOvnDnatRuleQueue, c.handleUpdateOvnDnatRule), time.Second, ctx.Done())
×
1545
        go wait.Until(runWorker("delete ovn dnat", c.delOvnDnatRuleQueue, c.handleDelOvnDnatRule), time.Second, ctx.Done())
×
1546

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

×
1549
        go wait.Until(runWorker("add ip", c.addIPQueue, c.handleAddReservedIP), time.Second, ctx.Done())
×
1550
        go wait.Until(runWorker("update ip", c.updateIPQueue, c.handleUpdateIP), time.Second, ctx.Done())
×
1551
        go wait.Until(runWorker("delete ip", c.delIPQueue, c.handleDelIP), time.Second, ctx.Done())
×
1552

×
1553
        go wait.Until(runWorker("add vip", c.addVirtualIPQueue, c.handleAddVirtualIP), time.Second, ctx.Done())
×
1554
        go wait.Until(runWorker("update vip", c.updateVirtualIPQueue, c.handleUpdateVirtualIP), time.Second, ctx.Done())
×
1555
        go wait.Until(runWorker("update virtual parent for vip", c.updateVirtualParentsQueue, c.handleUpdateVirtualParents), time.Second, ctx.Done())
×
1556
        go wait.Until(runWorker("delete vip", c.delVirtualIPQueue, c.handleDelVirtualIP), time.Second, ctx.Done())
×
1557

×
1558
        go wait.Until(runWorker("add iptables eip", c.addIptablesEipQueue, c.handleAddIptablesEip), time.Second, ctx.Done())
×
1559
        go wait.Until(runWorker("update iptables eip", c.updateIptablesEipQueue, c.handleUpdateIptablesEip), time.Second, ctx.Done())
×
1560
        go wait.Until(runWorker("reset iptables eip", c.resetIptablesEipQueue, c.handleResetIptablesEip), time.Second, ctx.Done())
×
1561
        go wait.Until(runWorker("delete iptables eip", c.delIptablesEipQueue, c.handleDelIptablesEip), time.Second, ctx.Done())
×
1562

×
1563
        go wait.Until(runWorker("add iptables fip", c.addIptablesFipQueue, c.handleAddIptablesFip), time.Second, ctx.Done())
×
1564
        go wait.Until(runWorker("update iptables fip", c.updateIptablesFipQueue, c.handleUpdateIptablesFip), time.Second, ctx.Done())
×
1565
        go wait.Until(runWorker("delete iptables fip", c.delIptablesFipQueue, c.handleDelIptablesFip), time.Second, ctx.Done())
×
1566

×
1567
        go wait.Until(runWorker("add iptables dnat rule", c.addIptablesDnatRuleQueue, c.handleAddIptablesDnatRule), time.Second, ctx.Done())
×
1568
        go wait.Until(runWorker("update iptables dnat rule", c.updateIptablesDnatRuleQueue, c.handleUpdateIptablesDnatRule), time.Second, ctx.Done())
×
1569
        go wait.Until(runWorker("delete iptables dnat rule", c.delIptablesDnatRuleQueue, c.handleDelIptablesDnatRule), time.Second, ctx.Done())
×
1570

×
1571
        go wait.Until(runWorker("add iptables snat rule", c.addIptablesSnatRuleQueue, c.handleAddIptablesSnatRule), time.Second, ctx.Done())
×
1572
        go wait.Until(runWorker("update iptables snat rule", c.updateIptablesSnatRuleQueue, c.handleUpdateIptablesSnatRule), time.Second, ctx.Done())
×
1573
        go wait.Until(runWorker("delete iptables snat rule", c.delIptablesSnatRuleQueue, c.handleDelIptablesSnatRule), time.Second, ctx.Done())
×
1574

×
1575
        go wait.Until(runWorker("add qos policy", c.addQoSPolicyQueue, c.handleAddQoSPolicy), time.Second, ctx.Done())
×
1576
        go wait.Until(runWorker("update qos policy", c.updateQoSPolicyQueue, c.handleUpdateQoSPolicy), time.Second, ctx.Done())
×
1577
        go wait.Until(runWorker("delete qos policy", c.delQoSPolicyQueue, c.handleDelQoSPolicy), time.Second, ctx.Done())
×
1578

×
1579
        if c.config.EnableANP {
×
1580
                go wait.Until(runWorker("add admin network policy", c.addAnpQueue, c.handleAddAnp), time.Second, ctx.Done())
×
1581
                go wait.Until(runWorker("update admin network policy", c.updateAnpQueue, c.handleUpdateAnp), time.Second, ctx.Done())
×
1582
                go wait.Until(runWorker("delete admin network policy", c.deleteAnpQueue, c.handleDeleteAnp), time.Second, ctx.Done())
×
1583

×
1584
                go wait.Until(runWorker("add base admin network policy", c.addBanpQueue, c.handleAddBanp), time.Second, ctx.Done())
×
1585
                go wait.Until(runWorker("update base admin network policy", c.updateBanpQueue, c.handleUpdateBanp), time.Second, ctx.Done())
×
1586
                go wait.Until(runWorker("delete base admin network policy", c.deleteBanpQueue, c.handleDeleteBanp), time.Second, ctx.Done())
×
1587

×
1588
                go wait.Until(runWorker("add cluster network policy", c.addCnpQueue, c.handleAddCnp), time.Second, ctx.Done())
×
1589
                go wait.Until(runWorker("update cluster network policy", c.updateCnpQueue, c.handleUpdateCnp), time.Second, ctx.Done())
×
1590
                go wait.Until(runWorker("delete cluster network policy", c.deleteCnpQueue, c.handleDeleteCnp), time.Second, ctx.Done())
×
1591
        }
×
1592

1593
        if c.config.EnableDNSNameResolver {
×
1594
                go wait.Until(runWorker("add or update dns name resolver", c.addOrUpdateDNSNameResolverQueue, c.handleAddOrUpdateDNSNameResolver), time.Second, ctx.Done())
×
1595
                go wait.Until(runWorker("delete dns name resolver", c.deleteDNSNameResolverQueue, c.handleDeleteDNSNameResolver), time.Second, ctx.Done())
×
1596
        }
×
1597

1598
        if c.config.EnableLiveMigrationOptimize {
×
1599
                go wait.Until(runWorker("add/update vmiMigration ", c.addOrUpdateVMIMigrationQueue, c.handleAddOrUpdateVMIMigration), 50*time.Millisecond, ctx.Done())
×
1600
        }
×
1601

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

×
1604
        go wait.Until(c.dbStatus, 15*time.Second, ctx.Done())
×
1605
}
1606

1607
func (c *Controller) allSubnetReady(subnets ...string) (bool, error) {
1✔
1608
        for _, lsName := range subnets {
2✔
1609
                exist, err := c.OVNNbClient.LogicalSwitchExists(lsName)
1✔
1610
                if err != nil {
1✔
1611
                        klog.Error(err)
×
1612
                        return false, fmt.Errorf("check logical switch %s exist: %w", lsName, err)
×
1613
                }
×
1614

1615
                if !exist {
2✔
1616
                        return false, nil
1✔
1617
                }
1✔
1618
        }
1619

1620
        return true, nil
1✔
1621
}
1622

1623
func (c *Controller) initResourceOnce() {
×
1624
        c.registerSubnetMetrics()
×
1625

×
1626
        if err := c.initNodeChassis(); err != nil {
×
1627
                util.LogFatalAndExit(err, "failed to initialize node chassis")
×
1628
        }
×
1629

1630
        if err := c.initDefaultDenyAllSecurityGroup(); err != nil {
×
1631
                util.LogFatalAndExit(err, "failed to initialize 'deny_all' security group")
×
1632
        }
×
1633
        if err := c.syncSecurityGroup(); err != nil {
×
1634
                util.LogFatalAndExit(err, "failed to sync security group")
×
1635
        }
×
1636

1637
        if err := c.syncVpcNatGatewayCR(); err != nil {
×
1638
                util.LogFatalAndExit(err, "failed to sync crd vpc nat gateways")
×
1639
        }
×
1640

1641
        if err := c.initVpcNatGw(); err != nil {
×
1642
                util.LogFatalAndExit(err, "failed to initialize vpc nat gateways")
×
1643
        }
×
1644
        if c.config.EnableLb {
×
1645
                if err := c.initVpcDNSConfig(); err != nil {
×
1646
                        util.LogFatalAndExit(err, "failed to initialize vpc-dns")
×
1647
                }
×
1648
        }
1649

1650
        // remove resources in ovndb that not exist any more in kubernetes resources
1651
        // process gc at last in case of affecting other init process
1652
        if err := c.gc(); err != nil {
×
1653
                util.LogFatalAndExit(err, "failed to run gc")
×
1654
        }
×
1655
}
1656

1657
func processNextWorkItem[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error, getItemKey func(any) string) bool {
×
1658
        item, shutdown := queue.Get()
×
1659
        if shutdown {
×
1660
                return false
×
1661
        }
×
1662

1663
        err := func(item T) error {
×
1664
                defer queue.Done(item)
×
1665
                if err := handler(item); err != nil {
×
1666
                        queue.AddRateLimited(item)
×
1667
                        return fmt.Errorf("error syncing %s %q: %w, requeuing", action, getItemKey(item), err)
×
1668
                }
×
1669
                queue.Forget(item)
×
1670
                return nil
×
1671
        }(item)
1672
        if err != nil {
×
1673
                utilruntime.HandleError(err)
×
1674
                return true
×
1675
        }
×
1676
        return true
×
1677
}
1678

1679
func getWorkItemKey(obj any) string {
×
1680
        switch v := obj.(type) {
×
1681
        case string:
×
1682
                return v
×
1683
        case *vpcService:
×
1684
                return cache.MetaObjectToName(obj.(*vpcService).Svc).String()
×
1685
        case *AdminNetworkPolicyChangedDelta:
×
1686
                return v.key
×
1687
        case *SwitchLBRuleInfo:
×
1688
                return v.Name
×
1689
        default:
×
1690
                key, err := cache.MetaNamespaceKeyFunc(obj)
×
1691
                if err != nil {
×
1692
                        utilruntime.HandleError(err)
×
1693
                        return ""
×
1694
                }
×
1695
                return key
×
1696
        }
1697
}
1698

1699
func runWorker[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error) func() {
×
1700
        return func() {
×
1701
                for processNextWorkItem(action, queue, handler, getWorkItemKey) {
×
1702
                }
×
1703
        }
1704
}
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