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

kubeovn / kube-ovn / 30331651685

28 Jul 2026 05:26AM UTC coverage: 25.905% (+0.08%) from 25.828%
30331651685

Pull #7067

github

changluyi
ci: run metallb e2e with three nodes

Signed-off-by: clyi <clyi@alauda.io>
Pull Request #7067: [release-1.16] fix(metallb): preserve internal underlay VIP traffic

113 of 282 new or added lines in 6 files covered. (40.07%)

275 existing lines in 3 files now uncovered.

14884 of 57455 relevant lines covered (25.91%)

0.3 hits per line

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

1.12
/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
        corev1 "k8s.io/api/core/v1"
17
        k8serrors "k8s.io/apimachinery/pkg/api/errors"
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
        "k8s.io/client-go/discovery"
23
        kubeinformers "k8s.io/client-go/informers"
24
        "k8s.io/client-go/kubernetes/scheme"
25
        typedcorev1 "k8s.io/client-go/kubernetes/typed/core/v1"
26
        appsv1 "k8s.io/client-go/listers/apps/v1"
27
        certListerv1 "k8s.io/client-go/listers/certificates/v1"
28
        v1 "k8s.io/client-go/listers/core/v1"
29
        discoveryv1 "k8s.io/client-go/listers/discovery/v1"
30
        netv1 "k8s.io/client-go/listers/networking/v1"
31
        "k8s.io/client-go/tools/cache"
32
        "k8s.io/client-go/tools/record"
33
        "k8s.io/client-go/util/workqueue"
34
        "k8s.io/klog/v2"
35
        "k8s.io/utils/keymutex"
36
        "k8s.io/utils/set"
37
        v1alpha1 "sigs.k8s.io/network-policy-api/apis/v1alpha1"
38
        netpolv1alpha2 "sigs.k8s.io/network-policy-api/apis/v1alpha2"
39
        anpinformer "sigs.k8s.io/network-policy-api/pkg/client/informers/externalversions"
40
        anplister "sigs.k8s.io/network-policy-api/pkg/client/listers/apis/v1alpha1"
41
        anplisterv1alpha2 "sigs.k8s.io/network-policy-api/pkg/client/listers/apis/v1alpha2"
42

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

52
const controllerAgentName = "kube-ovn-controller"
53

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

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

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

80
        OVNNbClient ovs.NbClient
81
        OVNSbClient ovs.SbClient
82

83
        // ExternalGatewayType define external gateway type, centralized
84
        ExternalGatewayType string
85

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

95
        vpcsLister           kubeovnlister.VpcLister
96
        vpcSynced            cache.InformerSynced
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
        switchLBRuleLister      kubeovnlister.SwitchLBRuleLister
132
        switchLBRuleSynced      cache.InformerSynced
133
        addSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[string]
134
        updateSwitchLBRuleQueue workqueue.TypedRateLimitingInterface[*SlrInfo]
135
        delSwitchLBRuleQueue    workqueue.TypedRateLimitingInterface[*SlrInfo]
136

137
        vpcDNSLister           kubeovnlister.VpcDnsLister
138
        vpcDNSSynced           cache.InformerSynced
139
        addOrUpdateVpcDNSQueue workqueue.TypedRateLimitingInterface[string]
140
        delVpcDNSQueue         workqueue.TypedRateLimitingInterface[string]
141

142
        subnetsLister           kubeovnlister.SubnetLister
143
        subnetSynced            cache.InformerSynced
144
        addOrUpdateSubnetQueue  workqueue.TypedRateLimitingInterface[string]
145
        deleteSubnetQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.Subnet]
146
        updateSubnetStatusQueue workqueue.TypedRateLimitingInterface[string]
147
        syncVirtualPortsQueue   workqueue.TypedRateLimitingInterface[string]
148
        subnetKeyMutex          keymutex.KeyMutex
149

150
        ippoolLister            kubeovnlister.IPPoolLister
151
        ippoolSynced            cache.InformerSynced
152
        addOrUpdateIPPoolQueue  workqueue.TypedRateLimitingInterface[string]
153
        updateIPPoolStatusQueue workqueue.TypedRateLimitingInterface[string]
154
        deleteIPPoolQueue       workqueue.TypedRateLimitingInterface[*kubeovnv1.IPPool]
155
        ippoolKeyMutex          keymutex.KeyMutex
156

157
        ipsLister     kubeovnlister.IPLister
158
        ipSynced      cache.InformerSynced
159
        addIPQueue    workqueue.TypedRateLimitingInterface[string]
160
        updateIPQueue workqueue.TypedRateLimitingInterface[string]
161
        delIPQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.IP]
162

163
        virtualIpsLister          kubeovnlister.VipLister
164
        virtualIpsSynced          cache.InformerSynced
165
        addVirtualIPQueue         workqueue.TypedRateLimitingInterface[string]
166
        updateVirtualIPQueue      workqueue.TypedRateLimitingInterface[string]
167
        updateVirtualParentsQueue workqueue.TypedRateLimitingInterface[string]
168
        delVirtualIPQueue         workqueue.TypedRateLimitingInterface[*kubeovnv1.Vip]
169

170
        iptablesEipsLister     kubeovnlister.IptablesEIPLister
171
        iptablesEipSynced      cache.InformerSynced
172
        addIptablesEipQueue    workqueue.TypedRateLimitingInterface[string]
173
        updateIptablesEipQueue workqueue.TypedRateLimitingInterface[string]
174
        resetIptablesEipQueue  workqueue.TypedRateLimitingInterface[string]
175
        delIptablesEipQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.IptablesEIP]
176

177
        iptablesFipsLister     kubeovnlister.IptablesFIPRuleLister
178
        iptablesFipSynced      cache.InformerSynced
179
        addIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
180
        updateIptablesFipQueue workqueue.TypedRateLimitingInterface[string]
181
        delIptablesFipQueue    workqueue.TypedRateLimitingInterface[string]
182

183
        iptablesDnatRulesLister     kubeovnlister.IptablesDnatRuleLister
184
        iptablesDnatRuleSynced      cache.InformerSynced
185
        addIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
186
        updateIptablesDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
187
        delIptablesDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
188

189
        iptablesSnatRulesLister     kubeovnlister.IptablesSnatRuleLister
190
        iptablesSnatRuleSynced      cache.InformerSynced
191
        addIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
192
        updateIptablesSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
193
        delIptablesSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
194

195
        ovnEipsLister     kubeovnlister.OvnEipLister
196
        ovnEipSynced      cache.InformerSynced
197
        addOvnEipQueue    workqueue.TypedRateLimitingInterface[string]
198
        updateOvnEipQueue workqueue.TypedRateLimitingInterface[string]
199
        resetOvnEipQueue  workqueue.TypedRateLimitingInterface[string]
200
        delOvnEipQueue    workqueue.TypedRateLimitingInterface[*kubeovnv1.OvnEip]
201

202
        ovnFipsLister     kubeovnlister.OvnFipLister
203
        ovnFipSynced      cache.InformerSynced
204
        addOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
205
        updateOvnFipQueue workqueue.TypedRateLimitingInterface[string]
206
        delOvnFipQueue    workqueue.TypedRateLimitingInterface[string]
207

208
        ovnSnatRulesLister     kubeovnlister.OvnSnatRuleLister
209
        ovnSnatRuleSynced      cache.InformerSynced
210
        addOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
211
        updateOvnSnatRuleQueue workqueue.TypedRateLimitingInterface[string]
212
        delOvnSnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
213

214
        ovnDnatRulesLister     kubeovnlister.OvnDnatRuleLister
215
        ovnDnatRuleSynced      cache.InformerSynced
216
        addOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
217
        updateOvnDnatRuleQueue workqueue.TypedRateLimitingInterface[string]
218
        delOvnDnatRuleQueue    workqueue.TypedRateLimitingInterface[string]
219

220
        providerNetworksLister kubeovnlister.ProviderNetworkLister
221
        providerNetworkSynced  cache.InformerSynced
222

223
        vlansLister     kubeovnlister.VlanLister
224
        vlanSynced      cache.InformerSynced
225
        addVlanQueue    workqueue.TypedRateLimitingInterface[string]
226
        delVlanQueue    workqueue.TypedRateLimitingInterface[string]
227
        updateVlanQueue workqueue.TypedRateLimitingInterface[string]
228
        vlanKeyMutex    keymutex.KeyMutex
229

230
        namespacesLister  v1.NamespaceLister
231
        namespacesSynced  cache.InformerSynced
232
        addNamespaceQueue workqueue.TypedRateLimitingInterface[string]
233
        nsKeyMutex        keymutex.KeyMutex
234

235
        nodesLister     v1.NodeLister
236
        nodesSynced     cache.InformerSynced
237
        addNodeQueue    workqueue.TypedRateLimitingInterface[string]
238
        updateNodeQueue workqueue.TypedRateLimitingInterface[string]
239
        deleteNodeQueue workqueue.TypedRateLimitingInterface[string]
240
        nodeKeyMutex    keymutex.KeyMutex
241

242
        servicesLister     v1.ServiceLister
243
        serviceSynced      cache.InformerSynced
244
        addServiceQueue    workqueue.TypedRateLimitingInterface[string]
245
        deleteServiceQueue workqueue.TypedRateLimitingInterface[*vpcService]
246
        updateServiceQueue workqueue.TypedRateLimitingInterface[*updateSvcObject]
247
        svcKeyMutex        keymutex.KeyMutex
248

249
        endpointSlicesLister          discoveryv1.EndpointSliceLister
250
        endpointSlicesSynced          cache.InformerSynced
251
        addOrUpdateEndpointSliceQueue workqueue.TypedRateLimitingInterface[string]
252
        epKeyMutex                    keymutex.KeyMutex
253
        serviceL2StatusMutex          sync.RWMutex
254
        serviceL2StatusIndexer        cache.Indexer
255
        serviceL2StatusSynced         cache.InformerSynced
256
        serviceL2StatusStarted        bool
257

258
        deploymentsLister appsv1.DeploymentLister
259
        deploymentsSynced cache.InformerSynced
260

261
        npsLister     netv1.NetworkPolicyLister
262
        npsSynced     cache.InformerSynced
263
        updateNpQueue workqueue.TypedRateLimitingInterface[string]
264
        deleteNpQueue workqueue.TypedRateLimitingInterface[string]
265
        npKeyMutex    keymutex.KeyMutex
266

267
        sgsLister          kubeovnlister.SecurityGroupLister
268
        sgSynced           cache.InformerSynced
269
        addOrUpdateSgQueue workqueue.TypedRateLimitingInterface[string]
270
        delSgQueue         workqueue.TypedRateLimitingInterface[string]
271
        syncSgPortsQueue   workqueue.TypedRateLimitingInterface[string]
272
        sgKeyMutex         keymutex.KeyMutex
273

274
        qosPoliciesLister    kubeovnlister.QoSPolicyLister
275
        qosPolicySynced      cache.InformerSynced
276
        addQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
277
        updateQoSPolicyQueue workqueue.TypedRateLimitingInterface[string]
278
        delQoSPolicyQueue    workqueue.TypedRateLimitingInterface[string]
279

280
        configMapsLister v1.ConfigMapLister
281
        configMapsSynced cache.InformerSynced
282

283
        anpsLister     anplister.AdminNetworkPolicyLister
284
        anpsSynced     cache.InformerSynced
285
        addAnpQueue    workqueue.TypedRateLimitingInterface[string]
286
        updateAnpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
287
        deleteAnpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.AdminNetworkPolicy]
288
        anpKeyMutex    keymutex.KeyMutex
289

290
        dnsNameResolversLister          kubeovnlister.DNSNameResolverLister
291
        dnsNameResolversSynced          cache.InformerSynced
292
        addOrUpdateDNSNameResolverQueue workqueue.TypedRateLimitingInterface[string]
293
        deleteDNSNameResolverQueue      workqueue.TypedRateLimitingInterface[*kubeovnv1.DNSNameResolver]
294

295
        banpsLister     anplister.BaselineAdminNetworkPolicyLister
296
        banpsSynced     cache.InformerSynced
297
        addBanpQueue    workqueue.TypedRateLimitingInterface[string]
298
        updateBanpQueue workqueue.TypedRateLimitingInterface[*AdminNetworkPolicyChangedDelta]
299
        deleteBanpQueue workqueue.TypedRateLimitingInterface[*v1alpha1.BaselineAdminNetworkPolicy]
300
        banpKeyMutex    keymutex.KeyMutex
301

302
        cnpsLister     anplisterv1alpha2.ClusterNetworkPolicyLister
303
        cnpsSynced     cache.InformerSynced
304
        addCnpQueue    workqueue.TypedRateLimitingInterface[string]
305
        updateCnpQueue workqueue.TypedRateLimitingInterface[*ClusterNetworkPolicyChangedDelta]
306
        deleteCnpQueue workqueue.TypedRateLimitingInterface[*netpolv1alpha2.ClusterNetworkPolicy]
307
        cnpKeyMutex    keymutex.KeyMutex
308

309
        csrLister           certListerv1.CertificateSigningRequestLister
310
        csrSynced           cache.InformerSynced
311
        addOrUpdateCsrQueue workqueue.TypedRateLimitingInterface[string]
312

313
        addOrUpdateVMIMigrationQueue workqueue.TypedRateLimitingInterface[string]
314
        deleteVMQueue                workqueue.TypedRateLimitingInterface[string]
315
        kubevirtInformerFactory      informer.KubeVirtInformerFactory
316

317
        netAttachLister          netAttachv1.NetworkAttachmentDefinitionLister
318
        netAttachSynced          cache.InformerSynced
319
        netAttachInformerFactory netAttach.SharedInformerFactory
320

321
        recorder               record.EventRecorder
322
        informerFactory        kubeinformers.SharedInformerFactory
323
        cmInformerFactory      kubeinformers.SharedInformerFactory
324
        deployInformerFactory  kubeinformers.SharedInformerFactory
325
        kubeovnInformerFactory kubeovninformer.SharedInformerFactory
326
        anpInformerFactory     anpinformer.SharedInformerFactory
327

328
        // Database health check
329
        dbFailureCount int
330

331
        distributedSubnetNeedSync atomic.Bool
332
}
333

334
func newTypedRateLimitingQueue[T comparable](name string, rateLimiter workqueue.TypedRateLimiter[T]) workqueue.TypedRateLimitingInterface[T] {
1✔
335
        if rateLimiter == nil {
2✔
336
                rateLimiter = workqueue.DefaultTypedControllerRateLimiter[T]()
1✔
337
        }
1✔
338
        return workqueue.NewTypedRateLimitingQueueWithConfig(rateLimiter, workqueue.TypedRateLimitingQueueConfig[T]{Name: name})
1✔
339
}
340

341
// Run creates and runs a new ovn controller
342
func Run(ctx context.Context, config *Configuration) {
×
343
        klog.V(4).Info("Creating event broadcaster")
×
344
        eventBroadcaster := record.NewBroadcasterWithCorrelatorOptions(record.CorrelatorOptions{BurstSize: 100})
×
345
        eventBroadcaster.StartLogging(klog.Infof)
×
346
        eventBroadcaster.StartRecordingToSink(&typedcorev1.EventSinkImpl{Interface: config.KubeFactoryClient.CoreV1().Events(metav1.NamespaceAll)})
×
347
        recorder := eventBroadcaster.NewRecorder(scheme.Scheme, corev1.EventSource{Component: controllerAgentName})
×
348
        custCrdRateLimiter := workqueue.NewTypedMaxOfRateLimiter(
×
349
                workqueue.NewTypedItemExponentialFailureRateLimiter[string](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
350
                &workqueue.TypedBucketRateLimiter[string]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
351
        )
×
352

×
353
        selector, err := labels.Parse(util.VpcEgressGatewayLabel)
×
354
        if err != nil {
×
355
                util.LogFatalAndExit(err, "failed to create label selector for vpc egress gateway workload")
×
356
        }
×
357

358
        informerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
359
                kubeinformers.WithTransform(util.TrimManagedFields),
×
360
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
361
                        listOption.AllowWatchBookmarks = true
×
362
                }))
×
363
        cmInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
364
                kubeinformers.WithNamespace(config.PodNamespace),
×
365
                kubeinformers.WithTransform(util.TrimManagedFields),
×
366
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
367
                        listOption.AllowWatchBookmarks = true
×
368
                }))
×
369
        // deployment informer used to list/watch vpc egress gateway workloads
370
        deployInformerFactory := kubeinformers.NewSharedInformerFactoryWithOptions(config.KubeFactoryClient, 0,
×
371
                kubeinformers.WithTransform(util.TrimManagedFields),
×
372
                kubeinformers.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
373
                        listOption.AllowWatchBookmarks = true
×
374
                        listOption.LabelSelector = selector.String()
×
375
                }))
×
376
        kubeovnInformerFactory := kubeovninformer.NewSharedInformerFactoryWithOptions(config.KubeOvnFactoryClient, 0,
×
377
                kubeovninformer.WithTransform(util.TrimManagedFields),
×
378
                kubeovninformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
379
                        listOption.AllowWatchBookmarks = true
×
380
                }))
×
381
        anpInformerFactory := anpinformer.NewSharedInformerFactoryWithOptions(config.AnpClient, 0,
×
382
                anpinformer.WithTransform(util.TrimManagedFields),
×
383
                anpinformer.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
384
                        listOption.AllowWatchBookmarks = true
×
385
                }))
×
386
        attachNetInformerFactory := netAttach.NewSharedInformerFactoryWithOptions(config.AttachNetClient, 0,
×
387
                netAttach.WithTweakListOptions(func(listOption *metav1.ListOptions) {
×
388
                        listOption.AllowWatchBookmarks = true
×
389
                }),
×
390
        )
391
        kubevirtInformerFactory := informer.NewKubeVirtInformerFactoryWithOptions(config.KubevirtClient.RestClient(), config.KubevirtClient,
×
392
                informer.WithTransform(util.TrimManagedFields),
×
393
        )
×
394

×
395
        vpcInformer := kubeovnInformerFactory.Kubeovn().V1().Vpcs()
×
396
        vpcNatGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcNatGateways()
×
397
        vpcEgressGatewayInformer := kubeovnInformerFactory.Kubeovn().V1().VpcEgressGateways()
×
398
        // BgpConf/EvpnConf informers are started lazily via StartBgpEvpnConfInformerFactory
×
399
        // because their CRDs are optional on clusters that don't use vpc-egress-gateway BGP/EVPN.
×
400
        subnetInformer := kubeovnInformerFactory.Kubeovn().V1().Subnets()
×
401
        ippoolInformer := kubeovnInformerFactory.Kubeovn().V1().IPPools()
×
402
        ipInformer := kubeovnInformerFactory.Kubeovn().V1().IPs()
×
403
        virtualIPInformer := kubeovnInformerFactory.Kubeovn().V1().Vips()
×
404
        iptablesEipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesEIPs()
×
405
        iptablesFipInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesFIPRules()
×
406
        iptablesDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesDnatRules()
×
407
        iptablesSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().IptablesSnatRules()
×
408
        vlanInformer := kubeovnInformerFactory.Kubeovn().V1().Vlans()
×
409
        providerNetworkInformer := kubeovnInformerFactory.Kubeovn().V1().ProviderNetworks()
×
410
        sgInformer := kubeovnInformerFactory.Kubeovn().V1().SecurityGroups()
×
411
        podInformer := informerFactory.Core().V1().Pods()
×
412
        namespaceInformer := informerFactory.Core().V1().Namespaces()
×
413
        nodeInformer := informerFactory.Core().V1().Nodes()
×
414
        serviceInformer := informerFactory.Core().V1().Services()
×
415
        endpointSliceInformer := informerFactory.Discovery().V1().EndpointSlices()
×
416
        deploymentInformer := deployInformerFactory.Apps().V1().Deployments()
×
417
        qosPolicyInformer := kubeovnInformerFactory.Kubeovn().V1().QoSPolicies()
×
418
        configMapInformer := cmInformerFactory.Core().V1().ConfigMaps()
×
419
        npInformer := informerFactory.Networking().V1().NetworkPolicies()
×
420
        switchLBRuleInformer := kubeovnInformerFactory.Kubeovn().V1().SwitchLBRules()
×
421
        vpcDNSInformer := kubeovnInformerFactory.Kubeovn().V1().VpcDnses()
×
422
        ovnEipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnEips()
×
423
        ovnFipInformer := kubeovnInformerFactory.Kubeovn().V1().OvnFips()
×
424
        ovnSnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnSnatRules()
×
425
        ovnDnatRuleInformer := kubeovnInformerFactory.Kubeovn().V1().OvnDnatRules()
×
426
        anpInformer := anpInformerFactory.Policy().V1alpha1().AdminNetworkPolicies()
×
427
        banpInformer := anpInformerFactory.Policy().V1alpha1().BaselineAdminNetworkPolicies()
×
428
        cnpInformer := anpInformerFactory.Policy().V1alpha2().ClusterNetworkPolicies()
×
429
        dnsNameResolverInformer := kubeovnInformerFactory.Kubeovn().V1().DNSNameResolvers()
×
430
        csrInformer := informerFactory.Certificates().V1().CertificateSigningRequests()
×
431
        netAttachInformer := attachNetInformerFactory.K8sCniCncfIo().V1().NetworkAttachmentDefinitions()
×
432

×
433
        numKeyLocks := max(runtime.NumCPU()*2, config.WorkerNum*2)
×
434
        controller := &Controller{
×
435
                config:             config,
×
436
                deletingPodObjMap:  xsync.NewMap[string, *corev1.Pod](),
×
437
                deletingNodeObjMap: xsync.NewMap[string, *corev1.Node](),
×
438
                ipam:               ovnipam.NewIPAM(),
×
439
                namedPort:          NewNamedPort(),
×
440

×
441
                vpcsLister:           vpcInformer.Lister(),
×
442
                vpcSynced:            vpcInformer.Informer().HasSynced,
×
443
                addOrUpdateVpcQueue:  newTypedRateLimitingQueue[string]("AddOrUpdateVpc", nil),
×
444
                vpcLastPoliciesMap:   xsync.NewMap[string, string](),
×
445
                delVpcQueue:          newTypedRateLimitingQueue[*kubeovnv1.Vpc]("DeleteVpc", nil),
×
446
                updateVpcStatusQueue: newTypedRateLimitingQueue[string]("UpdateVpcStatus", nil),
×
447
                vpcKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
448

×
449
                vpcNatGatewayLister:              vpcNatGatewayInformer.Lister(),
×
450
                vpcNatGatewaySynced:              vpcNatGatewayInformer.Informer().HasSynced,
×
451
                addOrUpdateVpcNatGatewayQueue:    newTypedRateLimitingQueue("AddOrUpdateVpcNatGw", custCrdRateLimiter),
×
452
                initVpcNatGatewayQueue:           newTypedRateLimitingQueue("InitVpcNatGw", custCrdRateLimiter),
×
453
                delVpcNatGatewayQueue:            newTypedRateLimitingQueue("DeleteVpcNatGw", custCrdRateLimiter),
×
454
                updateVpcEipQueue:                newTypedRateLimitingQueue("UpdateVpcEip", custCrdRateLimiter),
×
455
                updateVpcFloatingIPQueue:         newTypedRateLimitingQueue("UpdateVpcFloatingIp", custCrdRateLimiter),
×
456
                updateVpcDnatQueue:               newTypedRateLimitingQueue("UpdateVpcDnat", custCrdRateLimiter),
×
457
                updateVpcSnatQueue:               newTypedRateLimitingQueue("UpdateVpcSnat", custCrdRateLimiter),
×
458
                updateVpcSubnetQueue:             newTypedRateLimitingQueue("UpdateVpcSubnet", custCrdRateLimiter),
×
459
                vpcNatGwKeyMutex:                 keymutex.NewHashed(numKeyLocks),
×
460
                vpcNatGwExecKeyMutex:             keymutex.NewHashed(numKeyLocks),
×
461
                vpcEgressGatewayLister:           vpcEgressGatewayInformer.Lister(),
×
462
                vpcEgressGatewaySynced:           vpcEgressGatewayInformer.Informer().HasSynced,
×
463
                addOrUpdateVpcEgressGatewayQueue: newTypedRateLimitingQueue("AddOrUpdateVpcEgressGateway", custCrdRateLimiter),
×
464
                delVpcEgressGatewayQueue:         newTypedRateLimitingQueue("DeleteVpcEgressGateway", custCrdRateLimiter),
×
465
                vpcEgressGatewayKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
466

×
467
                // bgpConfLister/bgpConfSynced/evpnConfLister/evpnConfSynced are populated lazily
×
468
                // in startBgpEvpnConfInformer once the matching CRDs are detected.
×
469

×
470
                subnetsLister:           subnetInformer.Lister(),
×
471
                subnetSynced:            subnetInformer.Informer().HasSynced,
×
472
                addOrUpdateSubnetQueue:  newTypedRateLimitingQueue[string]("AddSubnet", nil),
×
473
                deleteSubnetQueue:       newTypedRateLimitingQueue[*kubeovnv1.Subnet]("DeleteSubnet", nil),
×
474
                updateSubnetStatusQueue: newTypedRateLimitingQueue[string]("UpdateSubnetStatus", nil),
×
475
                syncVirtualPortsQueue:   newTypedRateLimitingQueue[string]("SyncVirtualPort", nil),
×
476
                subnetKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
477

×
478
                ippoolLister:            ippoolInformer.Lister(),
×
479
                ippoolSynced:            ippoolInformer.Informer().HasSynced,
×
480
                addOrUpdateIPPoolQueue:  newTypedRateLimitingQueue[string]("AddIPPool", nil),
×
481
                updateIPPoolStatusQueue: newTypedRateLimitingQueue[string]("UpdateIPPoolStatus", nil),
×
482
                deleteIPPoolQueue:       newTypedRateLimitingQueue[*kubeovnv1.IPPool]("DeleteIPPool", nil),
×
483
                ippoolKeyMutex:          keymutex.NewHashed(numKeyLocks),
×
484

×
485
                ipsLister:     ipInformer.Lister(),
×
486
                ipSynced:      ipInformer.Informer().HasSynced,
×
487
                addIPQueue:    newTypedRateLimitingQueue[string]("AddIP", nil),
×
488
                updateIPQueue: newTypedRateLimitingQueue[string]("UpdateIP", nil),
×
489
                delIPQueue:    newTypedRateLimitingQueue[*kubeovnv1.IP]("DeleteIP", nil),
×
490

×
491
                virtualIpsLister:          virtualIPInformer.Lister(),
×
492
                virtualIpsSynced:          virtualIPInformer.Informer().HasSynced,
×
493
                addVirtualIPQueue:         newTypedRateLimitingQueue[string]("AddVirtualIP", nil),
×
494
                updateVirtualIPQueue:      newTypedRateLimitingQueue[string]("UpdateVirtualIP", nil),
×
495
                updateVirtualParentsQueue: newTypedRateLimitingQueue[string]("UpdateVirtualParents", nil),
×
496
                delVirtualIPQueue:         newTypedRateLimitingQueue[*kubeovnv1.Vip]("DeleteVirtualIP", nil),
×
497

×
498
                iptablesEipsLister:     iptablesEipInformer.Lister(),
×
499
                iptablesEipSynced:      iptablesEipInformer.Informer().HasSynced,
×
500
                addIptablesEipQueue:    newTypedRateLimitingQueue("AddIptablesEip", custCrdRateLimiter),
×
501
                updateIptablesEipQueue: newTypedRateLimitingQueue("UpdateIptablesEip", custCrdRateLimiter),
×
502
                resetIptablesEipQueue:  newTypedRateLimitingQueue("ResetIptablesEip", custCrdRateLimiter),
×
503
                delIptablesEipQueue:    newTypedRateLimitingQueue[*kubeovnv1.IptablesEIP]("DeleteIptablesEip", nil),
×
504

×
505
                iptablesFipsLister:     iptablesFipInformer.Lister(),
×
506
                iptablesFipSynced:      iptablesFipInformer.Informer().HasSynced,
×
507
                addIptablesFipQueue:    newTypedRateLimitingQueue("AddIptablesFip", custCrdRateLimiter),
×
508
                updateIptablesFipQueue: newTypedRateLimitingQueue("UpdateIptablesFip", custCrdRateLimiter),
×
509
                delIptablesFipQueue:    newTypedRateLimitingQueue("DeleteIptablesFip", custCrdRateLimiter),
×
510

×
511
                iptablesDnatRulesLister:     iptablesDnatRuleInformer.Lister(),
×
512
                iptablesDnatRuleSynced:      iptablesDnatRuleInformer.Informer().HasSynced,
×
513
                addIptablesDnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesDnatRule", custCrdRateLimiter),
×
514
                updateIptablesDnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesDnatRule", custCrdRateLimiter),
×
515
                delIptablesDnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesDnatRule", custCrdRateLimiter),
×
516

×
517
                iptablesSnatRulesLister:     iptablesSnatRuleInformer.Lister(),
×
518
                iptablesSnatRuleSynced:      iptablesSnatRuleInformer.Informer().HasSynced,
×
519
                addIptablesSnatRuleQueue:    newTypedRateLimitingQueue("AddIptablesSnatRule", custCrdRateLimiter),
×
520
                updateIptablesSnatRuleQueue: newTypedRateLimitingQueue("UpdateIptablesSnatRule", custCrdRateLimiter),
×
521
                delIptablesSnatRuleQueue:    newTypedRateLimitingQueue("DeleteIptablesSnatRule", custCrdRateLimiter),
×
522

×
523
                vlansLister:     vlanInformer.Lister(),
×
524
                vlanSynced:      vlanInformer.Informer().HasSynced,
×
525
                addVlanQueue:    newTypedRateLimitingQueue[string]("AddVlan", nil),
×
526
                delVlanQueue:    newTypedRateLimitingQueue[string]("DeleteVlan", nil),
×
527
                updateVlanQueue: newTypedRateLimitingQueue[string]("UpdateVlan", nil),
×
528
                vlanKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
529

×
530
                providerNetworksLister: providerNetworkInformer.Lister(),
×
531
                providerNetworkSynced:  providerNetworkInformer.Informer().HasSynced,
×
532

×
533
                podsLister:          podInformer.Lister(),
×
534
                podsSynced:          podInformer.Informer().HasSynced,
×
535
                addOrUpdatePodQueue: newTypedRateLimitingQueue[string]("AddOrUpdatePod", nil),
×
536
                deletePodQueue: workqueue.NewTypedRateLimitingQueueWithConfig(
×
537
                        workqueue.DefaultTypedControllerRateLimiter[string](),
×
538
                        workqueue.TypedRateLimitingQueueConfig[string]{
×
539
                                Name:          "DeletePod",
×
540
                                DelayingQueue: workqueue.NewTypedDelayingQueue[string](),
×
541
                        },
×
542
                ),
×
543
                updatePodSecurityQueue: newTypedRateLimitingQueue[string]("UpdatePodSecurity", nil),
×
544
                podKeyMutex:            keymutex.NewHashed(numKeyLocks),
×
545

×
546
                namespacesLister:  namespaceInformer.Lister(),
×
547
                namespacesSynced:  namespaceInformer.Informer().HasSynced,
×
548
                addNamespaceQueue: newTypedRateLimitingQueue[string]("AddNamespace", nil),
×
549
                nsKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
550

×
551
                nodesLister:     nodeInformer.Lister(),
×
552
                nodesSynced:     nodeInformer.Informer().HasSynced,
×
553
                addNodeQueue:    newTypedRateLimitingQueue[string]("AddNode", nil),
×
554
                updateNodeQueue: newTypedRateLimitingQueue[string]("UpdateNode", nil),
×
555
                deleteNodeQueue: newTypedRateLimitingQueue[string]("DeleteNode", nil),
×
556
                nodeKeyMutex:    keymutex.NewHashed(numKeyLocks),
×
557

×
558
                servicesLister:     serviceInformer.Lister(),
×
559
                serviceSynced:      serviceInformer.Informer().HasSynced,
×
560
                addServiceQueue:    newTypedRateLimitingQueue[string]("AddService", nil),
×
561
                deleteServiceQueue: newTypedRateLimitingQueue[*vpcService]("DeleteService", nil),
×
562
                updateServiceQueue: newTypedRateLimitingQueue[*updateSvcObject]("UpdateService", nil),
×
563
                svcKeyMutex:        keymutex.NewHashed(numKeyLocks),
×
564

×
565
                endpointSlicesLister:          endpointSliceInformer.Lister(),
×
566
                endpointSlicesSynced:          endpointSliceInformer.Informer().HasSynced,
×
567
                addOrUpdateEndpointSliceQueue: newTypedRateLimitingQueue[string]("UpdateEndpointSlice", nil),
×
568
                epKeyMutex:                    keymutex.NewHashed(numKeyLocks),
×
569

×
570
                deploymentsLister: deploymentInformer.Lister(),
×
571
                deploymentsSynced: deploymentInformer.Informer().HasSynced,
×
572

×
573
                qosPoliciesLister:    qosPolicyInformer.Lister(),
×
574
                qosPolicySynced:      qosPolicyInformer.Informer().HasSynced,
×
575
                addQoSPolicyQueue:    newTypedRateLimitingQueue("AddQoSPolicy", custCrdRateLimiter),
×
576
                updateQoSPolicyQueue: newTypedRateLimitingQueue("UpdateQoSPolicy", custCrdRateLimiter),
×
577
                delQoSPolicyQueue:    newTypedRateLimitingQueue("DeleteQoSPolicy", custCrdRateLimiter),
×
578

×
579
                configMapsLister: configMapInformer.Lister(),
×
580
                configMapsSynced: configMapInformer.Informer().HasSynced,
×
581

×
582
                sgKeyMutex:         keymutex.NewHashed(numKeyLocks),
×
583
                sgsLister:          sgInformer.Lister(),
×
584
                sgSynced:           sgInformer.Informer().HasSynced,
×
585
                addOrUpdateSgQueue: newTypedRateLimitingQueue[string]("UpdateSecurityGroup", nil),
×
586
                delSgQueue:         newTypedRateLimitingQueue[string]("DeleteSecurityGroup", nil),
×
587
                syncSgPortsQueue:   newTypedRateLimitingQueue[string]("SyncSecurityGroupPorts", nil),
×
588

×
589
                ovnEipsLister:     ovnEipInformer.Lister(),
×
590
                ovnEipSynced:      ovnEipInformer.Informer().HasSynced,
×
591
                addOvnEipQueue:    newTypedRateLimitingQueue("AddOvnEip", custCrdRateLimiter),
×
592
                updateOvnEipQueue: newTypedRateLimitingQueue("UpdateOvnEip", custCrdRateLimiter),
×
593
                resetOvnEipQueue:  newTypedRateLimitingQueue("ResetOvnEip", custCrdRateLimiter),
×
594
                delOvnEipQueue:    newTypedRateLimitingQueue[*kubeovnv1.OvnEip]("DeleteOvnEip", nil),
×
595

×
596
                ovnFipsLister:     ovnFipInformer.Lister(),
×
597
                ovnFipSynced:      ovnFipInformer.Informer().HasSynced,
×
598
                addOvnFipQueue:    newTypedRateLimitingQueue("AddOvnFip", custCrdRateLimiter),
×
599
                updateOvnFipQueue: newTypedRateLimitingQueue("UpdateOvnFip", custCrdRateLimiter),
×
600
                delOvnFipQueue:    newTypedRateLimitingQueue("DeleteOvnFip", custCrdRateLimiter),
×
601

×
602
                ovnSnatRulesLister:     ovnSnatRuleInformer.Lister(),
×
603
                ovnSnatRuleSynced:      ovnSnatRuleInformer.Informer().HasSynced,
×
604
                addOvnSnatRuleQueue:    newTypedRateLimitingQueue("AddOvnSnatRule", custCrdRateLimiter),
×
605
                updateOvnSnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnSnatRule", custCrdRateLimiter),
×
606
                delOvnSnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnSnatRule", custCrdRateLimiter),
×
607

×
608
                ovnDnatRulesLister:     ovnDnatRuleInformer.Lister(),
×
609
                ovnDnatRuleSynced:      ovnDnatRuleInformer.Informer().HasSynced,
×
610
                addOvnDnatRuleQueue:    newTypedRateLimitingQueue("AddOvnDnatRule", custCrdRateLimiter),
×
611
                updateOvnDnatRuleQueue: newTypedRateLimitingQueue("UpdateOvnDnatRule", custCrdRateLimiter),
×
612
                delOvnDnatRuleQueue:    newTypedRateLimitingQueue("DeleteOvnDnatRule", custCrdRateLimiter),
×
613

×
614
                csrLister:           csrInformer.Lister(),
×
615
                csrSynced:           csrInformer.Informer().HasSynced,
×
616
                addOrUpdateCsrQueue: newTypedRateLimitingQueue("AddOrUpdateCSR", custCrdRateLimiter),
×
617

×
618
                addOrUpdateVMIMigrationQueue: newTypedRateLimitingQueue[string]("AddOrUpdateVMIMigration", nil),
×
619
                deleteVMQueue:                newTypedRateLimitingQueue[string]("DeleteVM", nil),
×
620
                kubevirtInformerFactory:      kubevirtInformerFactory,
×
621

×
622
                netAttachLister:          netAttachInformer.Lister(),
×
623
                netAttachSynced:          netAttachInformer.Informer().HasSynced,
×
624
                netAttachInformerFactory: attachNetInformerFactory,
×
625

×
626
                recorder:               recorder,
×
627
                informerFactory:        informerFactory,
×
628
                cmInformerFactory:      cmInformerFactory,
×
629
                deployInformerFactory:  deployInformerFactory,
×
630
                kubeovnInformerFactory: kubeovnInformerFactory,
×
631
                anpInformerFactory:     anpInformerFactory,
×
632
        }
×
633

×
634
        if controller.OVNNbClient, err = ovs.NewOvnNbClient(
×
635
                config.OvnNbAddr,
×
636
                config.OvnTimeout,
×
637
                config.OvsDbConnectTimeout,
×
638
                config.OvsDbInactivityTimeout,
×
639
                config.OvsDbConnectMaxRetry,
×
640
        ); err != nil {
×
641
                util.LogFatalAndExit(err, "failed to create ovn nb client")
×
642
        }
×
643
        if controller.OVNSbClient, err = ovs.NewOvnSbClient(
×
644
                config.OvnSbAddr,
×
645
                config.OvnTimeout,
×
646
                config.OvsDbConnectTimeout,
×
647
                config.OvsDbInactivityTimeout,
×
648
                config.OvsDbConnectMaxRetry,
×
649
        ); err != nil {
×
650
                util.LogFatalAndExit(err, "failed to create ovn sb client")
×
651
        }
×
652
        if config.EnableLb {
×
653
                controller.switchLBRuleLister = switchLBRuleInformer.Lister()
×
654
                controller.switchLBRuleSynced = switchLBRuleInformer.Informer().HasSynced
×
655
                controller.addSwitchLBRuleQueue = newTypedRateLimitingQueue("AddSwitchLBRule", custCrdRateLimiter)
×
656
                controller.delSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
657
                        "DeleteSwitchLBRule",
×
658
                        workqueue.NewTypedMaxOfRateLimiter(
×
659
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SlrInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
660
                                &workqueue.TypedBucketRateLimiter[*SlrInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
661
                        ),
×
662
                )
×
663
                controller.updateSwitchLBRuleQueue = newTypedRateLimitingQueue(
×
664
                        "UpdateSwitchLBRule",
×
665
                        workqueue.NewTypedMaxOfRateLimiter(
×
666
                                workqueue.NewTypedItemExponentialFailureRateLimiter[*SlrInfo](time.Duration(config.CustCrdRetryMinDelay)*time.Second, time.Duration(config.CustCrdRetryMaxDelay)*time.Second),
×
667
                                &workqueue.TypedBucketRateLimiter[*SlrInfo]{Limiter: rate.NewLimiter(rate.Limit(10), 100)},
×
668
                        ),
×
669
                )
×
670

×
671
                controller.vpcDNSLister = vpcDNSInformer.Lister()
×
672
                controller.vpcDNSSynced = vpcDNSInformer.Informer().HasSynced
×
673
                controller.addOrUpdateVpcDNSQueue = newTypedRateLimitingQueue("AddOrUpdateVpcDns", custCrdRateLimiter)
×
674
                controller.delVpcDNSQueue = newTypedRateLimitingQueue("DeleteVpcDns", custCrdRateLimiter)
×
675
        }
×
676

677
        if config.EnableNP {
×
678
                controller.npsLister = npInformer.Lister()
×
679
                controller.npsSynced = npInformer.Informer().HasSynced
×
680
                controller.updateNpQueue = newTypedRateLimitingQueue[string]("UpdateNetworkPolicy", nil)
×
681
                controller.deleteNpQueue = newTypedRateLimitingQueue[string]("DeleteNetworkPolicy", nil)
×
682
                controller.npKeyMutex = keymutex.NewHashed(numKeyLocks)
×
683
        }
×
684

685
        if config.EnableANP {
×
686
                controller.anpsLister = anpInformer.Lister()
×
687
                controller.anpsSynced = anpInformer.Informer().HasSynced
×
688
                controller.addAnpQueue = newTypedRateLimitingQueue[string]("AddAdminNetworkPolicy", nil)
×
689
                controller.updateAnpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateAdminNetworkPolicy", nil)
×
690
                controller.deleteAnpQueue = newTypedRateLimitingQueue[*v1alpha1.AdminNetworkPolicy]("DeleteAdminNetworkPolicy", nil)
×
691
                controller.anpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
692

×
693
                controller.banpsLister = banpInformer.Lister()
×
694
                controller.banpsSynced = banpInformer.Informer().HasSynced
×
695
                controller.addBanpQueue = newTypedRateLimitingQueue[string]("AddBaseAdminNetworkPolicy", nil)
×
696
                controller.updateBanpQueue = newTypedRateLimitingQueue[*AdminNetworkPolicyChangedDelta]("UpdateBaseAdminNetworkPolicy", nil)
×
697
                controller.deleteBanpQueue = newTypedRateLimitingQueue[*v1alpha1.BaselineAdminNetworkPolicy]("DeleteBaseAdminNetworkPolicy", nil)
×
698
                controller.banpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
699

×
700
                controller.cnpsLister = cnpInformer.Lister()
×
701
                controller.cnpsSynced = cnpInformer.Informer().HasSynced
×
702
                controller.addCnpQueue = newTypedRateLimitingQueue[string]("AddClusterNetworkPolicy", nil)
×
703
                controller.updateCnpQueue = newTypedRateLimitingQueue[*ClusterNetworkPolicyChangedDelta]("UpdateClusterNetworkPolicy", nil)
×
704
                controller.deleteCnpQueue = newTypedRateLimitingQueue[*netpolv1alpha2.ClusterNetworkPolicy]("DeleteClusterNetworkPolicy", nil)
×
705
                controller.cnpKeyMutex = keymutex.NewHashed(numKeyLocks)
×
706
        }
×
707

708
        if config.EnableDNSNameResolver {
×
709
                if !config.EnableANP {
×
710
                        klog.Warning("DNS name resolver is enabled but ANP support is disabled, DNSNameResolver resources will not take effect")
×
711
                }
×
712
                controller.dnsNameResolversLister = dnsNameResolverInformer.Lister()
×
713
                controller.dnsNameResolversSynced = dnsNameResolverInformer.Informer().HasSynced
×
714
                controller.addOrUpdateDNSNameResolverQueue = newTypedRateLimitingQueue[string]("AddOrUpdateDNSNameResolver", nil)
×
715
                controller.deleteDNSNameResolverQueue = newTypedRateLimitingQueue[*kubeovnv1.DNSNameResolver]("DeleteDNSNameResolver", nil)
×
716
        }
717

718
        defer controller.shutdown()
×
719
        klog.Info("Starting OVN controller")
×
720

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

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

×
NEW
730
        // MetalLB is optional. When its ServiceL2Status API is available, use it to
×
NEW
731
        // identify the chassis announcing each underlay LoadBalancer VIP.
×
NEW
732
        controller.StartServiceL2StatusInformer(ctx)
×
NEW
733

×
734
        // Wait for the caches to be synced before starting workers
×
735
        controller.informerFactory.Start(ctx.Done())
×
736
        controller.cmInformerFactory.Start(ctx.Done())
×
737
        controller.deployInformerFactory.Start(ctx.Done())
×
738
        controller.kubeovnInformerFactory.Start(ctx.Done())
×
739
        controller.anpInformerFactory.Start(ctx.Done())
×
740
        controller.StartKubevirtInformerFactory(ctx, kubevirtInformerFactory)
×
741

×
742
        klog.Info("Waiting for informer caches to sync")
×
743
        cacheSyncs := []cache.InformerSynced{
×
744
                controller.vpcNatGatewaySynced, controller.vpcEgressGatewaySynced,
×
745
                controller.vpcSynced, controller.subnetSynced,
×
746
                controller.ipSynced, controller.virtualIpsSynced, controller.iptablesEipSynced,
×
747
                controller.iptablesFipSynced, controller.iptablesDnatRuleSynced, controller.iptablesSnatRuleSynced,
×
748
                controller.vlanSynced, controller.podsSynced, controller.namespacesSynced, controller.nodesSynced,
×
749
                controller.serviceSynced, controller.endpointSlicesSynced, controller.deploymentsSynced, controller.configMapsSynced,
×
750
                controller.ovnEipSynced, controller.ovnFipSynced, controller.ovnSnatRuleSynced,
×
751
                controller.ovnDnatRuleSynced,
×
752
        }
×
753
        if controller.config.EnableLb {
×
754
                cacheSyncs = append(cacheSyncs, controller.switchLBRuleSynced, controller.vpcDNSSynced)
×
755
        }
×
756
        if controller.config.EnableNP {
×
757
                cacheSyncs = append(cacheSyncs, controller.npsSynced)
×
758
        }
×
759
        if controller.config.EnableANP {
×
760
                cacheSyncs = append(cacheSyncs, controller.anpsSynced, controller.banpsSynced, controller.cnpsSynced)
×
761
        }
×
762
        if controller.config.EnableDNSNameResolver {
×
763
                cacheSyncs = append(cacheSyncs, controller.dnsNameResolversSynced)
×
764
        }
×
765

766
        if !cache.WaitForCacheSync(ctx.Done(), cacheSyncs...) {
×
767
                util.LogFatalAndExit(nil, "failed to wait for caches to sync")
×
768
        }
×
769

770
        if _, err = podInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
771
                AddFunc:    controller.enqueueAddPod,
×
772
                DeleteFunc: controller.enqueueDeletePod,
×
773
                UpdateFunc: controller.enqueueUpdatePod,
×
774
        }); err != nil {
×
775
                util.LogFatalAndExit(err, "failed to add pod event handler")
×
776
        }
×
777

778
        if _, err = namespaceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
779
                AddFunc:    controller.enqueueAddNamespace,
×
780
                UpdateFunc: controller.enqueueUpdateNamespace,
×
781
                DeleteFunc: controller.enqueueDeleteNamespace,
×
782
        }); err != nil {
×
783
                util.LogFatalAndExit(err, "failed to add namespace event handler")
×
784
        }
×
785

786
        if _, err = nodeInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
787
                AddFunc:    controller.enqueueAddNode,
×
788
                UpdateFunc: controller.enqueueUpdateNode,
×
789
                DeleteFunc: controller.enqueueDeleteNode,
×
790
        }); err != nil {
×
791
                util.LogFatalAndExit(err, "failed to add node event handler")
×
792
        }
×
793

794
        if _, err = serviceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
795
                AddFunc:    controller.enqueueAddService,
×
796
                DeleteFunc: controller.enqueueDeleteService,
×
797
                UpdateFunc: controller.enqueueUpdateService,
×
798
        }); err != nil {
×
799
                util.LogFatalAndExit(err, "failed to add service event handler")
×
800
        }
×
801

802
        if _, err = endpointSliceInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
803
                AddFunc:    controller.enqueueAddEndpointSlice,
×
804
                UpdateFunc: controller.enqueueUpdateEndpointSlice,
×
805
        }); err != nil {
×
806
                util.LogFatalAndExit(err, "failed to add endpoint slice event handler")
×
807
        }
×
808

809
        if _, err = deploymentInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
810
                AddFunc:    controller.enqueueAddDeployment,
×
811
                UpdateFunc: controller.enqueueUpdateDeployment,
×
812
        }); err != nil {
×
813
                util.LogFatalAndExit(err, "failed to add deployment event handler")
×
814
        }
×
815

816
        if _, err = vpcInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
817
                AddFunc:    controller.enqueueAddVpc,
×
818
                UpdateFunc: controller.enqueueUpdateVpc,
×
819
                DeleteFunc: controller.enqueueDelVpc,
×
820
        }); err != nil {
×
821
                util.LogFatalAndExit(err, "failed to add vpc event handler")
×
822
        }
×
823

824
        if _, err = vpcNatGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
825
                AddFunc:    controller.enqueueAddVpcNatGw,
×
826
                UpdateFunc: controller.enqueueUpdateVpcNatGw,
×
827
                DeleteFunc: controller.enqueueDeleteVpcNatGw,
×
828
        }); err != nil {
×
829
                util.LogFatalAndExit(err, "failed to add vpc nat gateway event handler")
×
830
        }
×
831

832
        if _, err = vpcEgressGatewayInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
833
                AddFunc:    controller.enqueueAddVpcEgressGateway,
×
834
                UpdateFunc: controller.enqueueUpdateVpcEgressGateway,
×
835
                DeleteFunc: controller.enqueueDeleteVpcEgressGateway,
×
836
        }); err != nil {
×
837
                util.LogFatalAndExit(err, "failed to add vpc egress gateway event handler")
×
838
        }
×
839

840
        if _, err = subnetInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
841
                AddFunc:    controller.enqueueAddSubnet,
×
842
                UpdateFunc: controller.enqueueUpdateSubnet,
×
843
                DeleteFunc: controller.enqueueDeleteSubnet,
×
844
        }); err != nil {
×
845
                util.LogFatalAndExit(err, "failed to add subnet event handler")
×
846
        }
×
847

848
        if _, err = ippoolInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
849
                AddFunc:    controller.enqueueAddIPPool,
×
850
                UpdateFunc: controller.enqueueUpdateIPPool,
×
851
                DeleteFunc: controller.enqueueDeleteIPPool,
×
852
        }); err != nil {
×
853
                util.LogFatalAndExit(err, "failed to add ippool event handler")
×
854
        }
×
855

856
        if _, err = ipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
857
                AddFunc:    controller.enqueueAddIP,
×
858
                UpdateFunc: controller.enqueueUpdateIP,
×
859
                DeleteFunc: controller.enqueueDelIP,
×
860
        }); err != nil {
×
861
                util.LogFatalAndExit(err, "failed to add ips event handler")
×
862
        }
×
863

864
        if _, err = vlanInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
865
                AddFunc:    controller.enqueueAddVlan,
×
866
                DeleteFunc: controller.enqueueDelVlan,
×
867
                UpdateFunc: controller.enqueueUpdateVlan,
×
868
        }); err != nil {
×
869
                util.LogFatalAndExit(err, "failed to add vlan event handler")
×
870
        }
×
871

872
        if _, err = sgInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
873
                AddFunc:    controller.enqueueAddSg,
×
874
                DeleteFunc: controller.enqueueDeleteSg,
×
875
                UpdateFunc: controller.enqueueUpdateSg,
×
876
        }); err != nil {
×
877
                util.LogFatalAndExit(err, "failed to add security group event handler")
×
878
        }
×
879

880
        if _, err = virtualIPInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
881
                AddFunc:    controller.enqueueAddVirtualIP,
×
882
                UpdateFunc: controller.enqueueUpdateVirtualIP,
×
883
                DeleteFunc: controller.enqueueDelVirtualIP,
×
884
        }); err != nil {
×
885
                util.LogFatalAndExit(err, "failed to add virtual ip event handler")
×
886
        }
×
887

888
        if _, err = iptablesEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
889
                AddFunc:    controller.enqueueAddIptablesEip,
×
890
                UpdateFunc: controller.enqueueUpdateIptablesEip,
×
891
                DeleteFunc: controller.enqueueDelIptablesEip,
×
892
        }); err != nil {
×
893
                util.LogFatalAndExit(err, "failed to add iptables eip event handler")
×
894
        }
×
895

896
        if _, err = iptablesFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
897
                AddFunc:    controller.enqueueAddIptablesFip,
×
898
                UpdateFunc: controller.enqueueUpdateIptablesFip,
×
899
                DeleteFunc: controller.enqueueDelIptablesFip,
×
900
        }); err != nil {
×
901
                util.LogFatalAndExit(err, "failed to add iptables fip event handler")
×
902
        }
×
903

904
        if _, err = iptablesDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
905
                AddFunc:    controller.enqueueAddIptablesDnatRule,
×
906
                UpdateFunc: controller.enqueueUpdateIptablesDnatRule,
×
907
                DeleteFunc: controller.enqueueDelIptablesDnatRule,
×
908
        }); err != nil {
×
909
                util.LogFatalAndExit(err, "failed to add iptables dnat event handler")
×
910
        }
×
911

912
        if _, err = iptablesSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
913
                AddFunc:    controller.enqueueAddIptablesSnatRule,
×
914
                UpdateFunc: controller.enqueueUpdateIptablesSnatRule,
×
915
                DeleteFunc: controller.enqueueDelIptablesSnatRule,
×
916
        }); err != nil {
×
917
                util.LogFatalAndExit(err, "failed to add iptables snat rule event handler")
×
918
        }
×
919

920
        if _, err = ovnEipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
921
                AddFunc:    controller.enqueueAddOvnEip,
×
922
                UpdateFunc: controller.enqueueUpdateOvnEip,
×
923
                DeleteFunc: controller.enqueueDelOvnEip,
×
924
        }); err != nil {
×
925
                util.LogFatalAndExit(err, "failed to add ovn eip event handler")
×
926
        }
×
927

928
        if _, err = ovnFipInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
929
                AddFunc:    controller.enqueueAddOvnFip,
×
930
                UpdateFunc: controller.enqueueUpdateOvnFip,
×
931
                DeleteFunc: controller.enqueueDelOvnFip,
×
932
        }); err != nil {
×
933
                util.LogFatalAndExit(err, "failed to add ovn fip event handler")
×
934
        }
×
935

936
        if _, err = ovnSnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
937
                AddFunc:    controller.enqueueAddOvnSnatRule,
×
938
                UpdateFunc: controller.enqueueUpdateOvnSnatRule,
×
939
                DeleteFunc: controller.enqueueDelOvnSnatRule,
×
940
        }); err != nil {
×
941
                util.LogFatalAndExit(err, "failed to add ovn snat rule event handler")
×
942
        }
×
943

944
        if _, err = ovnDnatRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
945
                AddFunc:    controller.enqueueAddOvnDnatRule,
×
946
                UpdateFunc: controller.enqueueUpdateOvnDnatRule,
×
947
                DeleteFunc: controller.enqueueDelOvnDnatRule,
×
948
        }); err != nil {
×
949
                util.LogFatalAndExit(err, "failed to add ovn dnat rule event handler")
×
950
        }
×
951

952
        if _, err = qosPolicyInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
953
                AddFunc:    controller.enqueueAddQoSPolicy,
×
954
                UpdateFunc: controller.enqueueUpdateQoSPolicy,
×
955
                DeleteFunc: controller.enqueueDelQoSPolicy,
×
956
        }); err != nil {
×
957
                util.LogFatalAndExit(err, "failed to add qos policy event handler")
×
958
        }
×
959

960
        if config.EnableLb {
×
961
                if _, err = switchLBRuleInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
962
                        AddFunc:    controller.enqueueAddSwitchLBRule,
×
963
                        UpdateFunc: controller.enqueueUpdateSwitchLBRule,
×
964
                        DeleteFunc: controller.enqueueDeleteSwitchLBRule,
×
965
                }); err != nil {
×
966
                        util.LogFatalAndExit(err, "failed to add switch lb rule event handler")
×
967
                }
×
968

969
                if _, err = vpcDNSInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
970
                        AddFunc:    controller.enqueueAddVpcDNS,
×
971
                        UpdateFunc: controller.enqueueUpdateVpcDNS,
×
972
                        DeleteFunc: controller.enqueueDeleteVPCDNS,
×
973
                }); err != nil {
×
974
                        util.LogFatalAndExit(err, "failed to add vpc dns event handler")
×
975
                }
×
976
        }
977

978
        if config.EnableNP {
×
979
                if _, err = npInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
980
                        AddFunc:    controller.enqueueAddNp,
×
981
                        UpdateFunc: controller.enqueueUpdateNp,
×
982
                        DeleteFunc: controller.enqueueDeleteNp,
×
983
                }); err != nil {
×
984
                        util.LogFatalAndExit(err, "failed to add network policy event handler")
×
985
                }
×
986
        }
987

988
        if config.EnableANP {
×
989
                if _, err = anpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
990
                        AddFunc:    controller.enqueueAddAnp,
×
991
                        UpdateFunc: controller.enqueueUpdateAnp,
×
992
                        DeleteFunc: controller.enqueueDeleteAnp,
×
993
                }); err != nil {
×
994
                        util.LogFatalAndExit(err, "failed to add admin network policy event handler")
×
995
                }
×
996

997
                if _, err = banpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
998
                        AddFunc:    controller.enqueueAddBanp,
×
999
                        UpdateFunc: controller.enqueueUpdateBanp,
×
1000
                        DeleteFunc: controller.enqueueDeleteBanp,
×
1001
                }); err != nil {
×
1002
                        util.LogFatalAndExit(err, "failed to add baseline admin network policy event handler")
×
1003
                }
×
1004

1005
                if _, err = cnpInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1006
                        AddFunc:    controller.enqueueAddCnp,
×
1007
                        UpdateFunc: controller.enqueueUpdateCnp,
×
1008
                        DeleteFunc: controller.enqueueDeleteCnp,
×
1009
                }); err != nil {
×
1010
                        util.LogFatalAndExit(err, "failed to add cluster network policy event handler")
×
1011
                }
×
1012

1013
                maxPriorityPerMap := util.CnpMaxPriority + 1
×
1014
                controller.anpPrioNameMap = make(map[int32]string, maxPriorityPerMap)
×
1015
                controller.anpNamePrioMap = make(map[string]int32, maxPriorityPerMap)
×
1016
                controller.bnpPrioNameMap = make(map[int32]string, maxPriorityPerMap)
×
1017
                controller.bnpNamePrioMap = make(map[string]int32, maxPriorityPerMap)
×
1018
        }
1019

1020
        if config.EnableDNSNameResolver {
×
1021
                if _, err = dnsNameResolverInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1022
                        AddFunc:    controller.enqueueAddDNSNameResolver,
×
1023
                        UpdateFunc: controller.enqueueUpdateDNSNameResolver,
×
1024
                        DeleteFunc: controller.enqueueDeleteDNSNameResolver,
×
1025
                }); err != nil {
×
1026
                        util.LogFatalAndExit(err, "failed to add dns name resolver event handler")
×
1027
                }
×
1028
        }
1029

1030
        if config.EnableOVNIPSec {
×
1031
                if _, err = csrInformer.Informer().AddEventHandler(cache.ResourceEventHandlerFuncs{
×
1032
                        AddFunc:    controller.enqueueAddCsr,
×
1033
                        UpdateFunc: controller.enqueueUpdateCsr,
×
1034
                        // no need to add delete func for csr
×
1035
                }); err != nil {
×
1036
                        util.LogFatalAndExit(err, "failed to add csr event handler")
×
1037
                }
×
1038
        }
1039

1040
        controller.Run(ctx)
×
1041
}
1042

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

1054
        if err := c.OVNNbClient.SetUseCtInvMatch(); err != nil {
×
1055
                util.LogFatalAndExit(err, "failed to set NB_Global option use_ct_inv_match to false")
×
1056
        }
×
1057

1058
        if err := c.OVNNbClient.SetLsCtSkipDstLportIPs(c.config.LsCtSkipDstLportIPs); err != nil {
×
1059
                util.LogFatalAndExit(err, "failed to set NB_Global option ls_ct_skip_dst_lport_ips")
×
1060
        }
×
1061

1062
        if err := c.OVNNbClient.SetNodeLocalDNSIP(strings.Join(c.config.NodeLocalDNSIPs, ",")); err != nil {
×
1063
                util.LogFatalAndExit(err, "failed to set NB_Global option node_local_dns_ip")
×
1064
        }
×
1065

1066
        if err := c.OVNNbClient.SetSkipConntrackCidrs(c.config.SkipConntrackDstCidrs); err != nil {
×
1067
                util.LogFatalAndExit(err, "failed to set NB_Global option skip_conntrack_ipcidrs")
×
1068
        }
×
1069

1070
        if err := c.OVNNbClient.SetOVNIPSec(c.config.EnableOVNIPSec); err != nil {
×
1071
                util.LogFatalAndExit(err, "failed to set NB_Global ipsec")
×
1072
        }
×
1073

1074
        if err := c.InitOVN(); err != nil {
×
1075
                util.LogFatalAndExit(err, "failed to initialize ovn resources")
×
1076
        }
×
1077

1078
        // sync ip crd before initIPAM since ip crd will be used to restore vm and statefulset pod in initIPAM
1079
        if err := c.syncIPCR(); err != nil {
×
1080
                util.LogFatalAndExit(err, "failed to sync crd ips")
×
1081
        }
×
1082

1083
        if err := c.syncFinalizers(); err != nil {
×
1084
                util.LogFatalAndExit(err, "failed to initialize crd finalizers")
×
1085
        }
×
1086

1087
        if err := c.InitIPAM(); err != nil {
×
1088
                util.LogFatalAndExit(err, "failed to initialize ipam")
×
1089
        }
×
1090

1091
        if err := c.syncNodeRoutes(); err != nil {
×
1092
                util.LogFatalAndExit(err, "failed to initialize node routes")
×
1093
        }
×
1094

1095
        if err := c.syncSubnetCR(); err != nil {
×
1096
                util.LogFatalAndExit(err, "failed to sync crd subnets")
×
1097
        }
×
1098

1099
        if err := c.syncVlanCR(); err != nil {
×
1100
                util.LogFatalAndExit(err, "failed to sync crd vlans")
×
1101
        }
×
1102

1103
        if c.config.EnableOVNIPSec && !c.config.CertManagerIPSecCert {
×
1104
                if err := c.InitDefaultOVNIPsecCA(); err != nil {
×
1105
                        util.LogFatalAndExit(err, "failed to init ovn ipsec CA")
×
1106
                }
×
1107
        }
1108

1109
        c.startKubeOVNTLSManager(ctx)
×
1110

×
1111
        // start workers to do all the network operations
×
1112
        c.startWorkers(ctx)
×
1113

×
1114
        c.initResourceOnce()
×
1115
        <-ctx.Done()
×
1116
        klog.Info("Shutting down workers")
×
1117

×
1118
        c.OVNNbClient.Close()
×
1119
        c.OVNSbClient.Close()
×
1120
}
1121

1122
func (c *Controller) dbStatus() {
×
1123
        const maxFailures = 5
×
1124

×
1125
        done := make(chan error, 2)
×
1126
        go func() {
×
1127
                done <- c.OVNNbClient.Echo(context.Background())
×
1128
        }()
×
1129
        go func() {
×
1130
                done <- c.OVNSbClient.Echo(context.Background())
×
1131
        }()
×
1132

1133
        resultsReceived := 0
×
1134
        timeout := time.After(time.Duration(c.config.OvnTimeout) * time.Second)
×
1135

×
1136
        for resultsReceived < 2 {
×
1137
                select {
×
1138
                case err := <-done:
×
1139
                        resultsReceived++
×
1140
                        if err != nil {
×
1141
                                c.dbFailureCount++
×
1142
                                klog.Errorf("OVN database echo failed (%d/%d): %v", c.dbFailureCount, maxFailures, err)
×
1143
                                if c.dbFailureCount >= maxFailures {
×
1144
                                        util.LogFatalAndExit(err, "OVN database connection failed after %d attempts", maxFailures)
×
1145
                                }
×
1146
                                return
×
1147
                        }
1148
                case <-timeout:
×
1149
                        c.dbFailureCount++
×
1150
                        klog.Errorf("OVN database echo timeout (%d/%d) after %ds", c.dbFailureCount, maxFailures, c.config.OvnTimeout)
×
1151
                        if c.dbFailureCount >= maxFailures {
×
1152
                                util.LogFatalAndExit(nil, "OVN database connection timeout after %d attempts", maxFailures)
×
1153
                        }
×
1154
                        return
×
1155
                }
1156
        }
1157

1158
        if c.dbFailureCount > 0 {
×
1159
                klog.Infof("OVN database connection recovered after %d failures", c.dbFailureCount)
×
1160
                c.dbFailureCount = 0
×
1161
        }
×
1162
}
1163

1164
func (c *Controller) shutdown() {
×
1165
        utilruntime.HandleCrash()
×
1166

×
1167
        c.addOrUpdatePodQueue.ShutDown()
×
1168
        c.deletePodQueue.ShutDown()
×
1169
        c.updatePodSecurityQueue.ShutDown()
×
1170

×
1171
        c.addNamespaceQueue.ShutDown()
×
1172

×
1173
        c.addOrUpdateSubnetQueue.ShutDown()
×
1174
        c.deleteSubnetQueue.ShutDown()
×
1175
        c.updateSubnetStatusQueue.ShutDown()
×
1176
        c.syncVirtualPortsQueue.ShutDown()
×
1177

×
1178
        c.addOrUpdateIPPoolQueue.ShutDown()
×
1179
        c.updateIPPoolStatusQueue.ShutDown()
×
1180
        c.deleteIPPoolQueue.ShutDown()
×
1181

×
1182
        c.addNodeQueue.ShutDown()
×
1183
        c.updateNodeQueue.ShutDown()
×
1184
        c.deleteNodeQueue.ShutDown()
×
1185

×
1186
        c.addServiceQueue.ShutDown()
×
1187
        c.deleteServiceQueue.ShutDown()
×
1188
        c.updateServiceQueue.ShutDown()
×
1189
        c.addOrUpdateEndpointSliceQueue.ShutDown()
×
1190

×
1191
        c.addVlanQueue.ShutDown()
×
1192
        c.delVlanQueue.ShutDown()
×
1193
        c.updateVlanQueue.ShutDown()
×
1194

×
1195
        c.addOrUpdateVpcQueue.ShutDown()
×
1196
        c.updateVpcStatusQueue.ShutDown()
×
1197
        c.delVpcQueue.ShutDown()
×
1198

×
1199
        c.addOrUpdateVpcNatGatewayQueue.ShutDown()
×
1200
        c.initVpcNatGatewayQueue.ShutDown()
×
1201
        c.delVpcNatGatewayQueue.ShutDown()
×
1202
        c.updateVpcEipQueue.ShutDown()
×
1203
        c.updateVpcFloatingIPQueue.ShutDown()
×
1204
        c.updateVpcDnatQueue.ShutDown()
×
1205
        c.updateVpcSnatQueue.ShutDown()
×
1206
        c.updateVpcSubnetQueue.ShutDown()
×
1207

×
1208
        c.addOrUpdateVpcEgressGatewayQueue.ShutDown()
×
1209
        c.delVpcEgressGatewayQueue.ShutDown()
×
1210

×
1211
        if c.config.EnableLb {
×
1212
                c.addSwitchLBRuleQueue.ShutDown()
×
1213
                c.delSwitchLBRuleQueue.ShutDown()
×
1214
                c.updateSwitchLBRuleQueue.ShutDown()
×
1215

×
1216
                c.addOrUpdateVpcDNSQueue.ShutDown()
×
1217
                c.delVpcDNSQueue.ShutDown()
×
1218
        }
×
1219

1220
        c.addIPQueue.ShutDown()
×
1221
        c.updateIPQueue.ShutDown()
×
1222
        c.delIPQueue.ShutDown()
×
1223

×
1224
        c.addVirtualIPQueue.ShutDown()
×
1225
        c.updateVirtualIPQueue.ShutDown()
×
1226
        c.updateVirtualParentsQueue.ShutDown()
×
1227
        c.delVirtualIPQueue.ShutDown()
×
1228

×
1229
        c.addIptablesEipQueue.ShutDown()
×
1230
        c.updateIptablesEipQueue.ShutDown()
×
1231
        c.resetIptablesEipQueue.ShutDown()
×
1232
        c.delIptablesEipQueue.ShutDown()
×
1233

×
1234
        c.addIptablesFipQueue.ShutDown()
×
1235
        c.updateIptablesFipQueue.ShutDown()
×
1236
        c.delIptablesFipQueue.ShutDown()
×
1237

×
1238
        c.addIptablesDnatRuleQueue.ShutDown()
×
1239
        c.updateIptablesDnatRuleQueue.ShutDown()
×
1240
        c.delIptablesDnatRuleQueue.ShutDown()
×
1241

×
1242
        c.addIptablesSnatRuleQueue.ShutDown()
×
1243
        c.updateIptablesSnatRuleQueue.ShutDown()
×
1244
        c.delIptablesSnatRuleQueue.ShutDown()
×
1245

×
1246
        c.addQoSPolicyQueue.ShutDown()
×
1247
        c.updateQoSPolicyQueue.ShutDown()
×
1248
        c.delQoSPolicyQueue.ShutDown()
×
1249

×
1250
        c.addOvnEipQueue.ShutDown()
×
1251
        c.updateOvnEipQueue.ShutDown()
×
1252
        c.resetOvnEipQueue.ShutDown()
×
1253
        c.delOvnEipQueue.ShutDown()
×
1254

×
1255
        c.addOvnFipQueue.ShutDown()
×
1256
        c.updateOvnFipQueue.ShutDown()
×
1257
        c.delOvnFipQueue.ShutDown()
×
1258

×
1259
        c.addOvnSnatRuleQueue.ShutDown()
×
1260
        c.updateOvnSnatRuleQueue.ShutDown()
×
1261
        c.delOvnSnatRuleQueue.ShutDown()
×
1262

×
1263
        c.addOvnDnatRuleQueue.ShutDown()
×
1264
        c.updateOvnDnatRuleQueue.ShutDown()
×
1265
        c.delOvnDnatRuleQueue.ShutDown()
×
1266

×
1267
        if c.config.EnableNP {
×
1268
                c.updateNpQueue.ShutDown()
×
1269
                c.deleteNpQueue.ShutDown()
×
1270
        }
×
1271
        if c.config.EnableANP {
×
1272
                c.addAnpQueue.ShutDown()
×
1273
                c.updateAnpQueue.ShutDown()
×
1274
                c.deleteAnpQueue.ShutDown()
×
1275

×
1276
                c.addBanpQueue.ShutDown()
×
1277
                c.updateBanpQueue.ShutDown()
×
1278
                c.deleteBanpQueue.ShutDown()
×
1279

×
1280
                c.addCnpQueue.ShutDown()
×
1281
                c.updateCnpQueue.ShutDown()
×
1282
                c.deleteCnpQueue.ShutDown()
×
1283
        }
×
1284

1285
        if c.config.EnableDNSNameResolver {
×
1286
                c.addOrUpdateDNSNameResolverQueue.ShutDown()
×
1287
                c.deleteDNSNameResolverQueue.ShutDown()
×
1288
        }
×
1289

1290
        c.addOrUpdateSgQueue.ShutDown()
×
1291
        c.delSgQueue.ShutDown()
×
1292
        c.syncSgPortsQueue.ShutDown()
×
1293

×
1294
        c.addOrUpdateCsrQueue.ShutDown()
×
1295

×
1296
        if c.config.EnableLiveMigrationOptimize {
×
1297
                c.addOrUpdateVMIMigrationQueue.ShutDown()
×
1298
        }
×
1299
}
1300

1301
func (c *Controller) startWorkers(ctx context.Context) {
×
1302
        klog.Info("Starting workers")
×
1303

×
1304
        go wait.Until(runWorker("add/update vpc", c.addOrUpdateVpcQueue, c.handleAddOrUpdateVpc), time.Second, ctx.Done())
×
1305
        go wait.Until(runWorker("delete vpc", c.delVpcQueue, c.handleDelVpc), time.Second, ctx.Done())
×
1306
        go wait.Until(runWorker("update status of vpc", c.updateVpcStatusQueue, c.handleUpdateVpcStatus), time.Second, ctx.Done())
×
1307

×
1308
        go wait.Until(runWorker("add/update vpc nat gateway", c.addOrUpdateVpcNatGatewayQueue, c.handleAddOrUpdateVpcNatGw), time.Second, ctx.Done())
×
1309
        go wait.Until(runWorker("init vpc nat gateway", c.initVpcNatGatewayQueue, c.handleInitVpcNatGw), time.Second, ctx.Done())
×
1310
        go wait.Until(runWorker("delete vpc nat gateway", c.delVpcNatGatewayQueue, c.handleDelVpcNatGw), time.Second, ctx.Done())
×
1311
        go wait.Until(runWorker("add/update vpc egress gateway", c.addOrUpdateVpcEgressGatewayQueue, c.handleAddOrUpdateVpcEgressGateway), time.Second, ctx.Done())
×
1312
        go wait.Until(runWorker("delete vpc egress gateway", c.delVpcEgressGatewayQueue, c.handleDelVpcEgressGateway), time.Second, ctx.Done())
×
1313
        go wait.Until(runWorker("update fip for vpc nat gateway", c.updateVpcFloatingIPQueue, c.handleUpdateVpcFloatingIP), time.Second, ctx.Done())
×
1314
        go wait.Until(runWorker("update eip for vpc nat gateway", c.updateVpcEipQueue, c.handleUpdateVpcEip), time.Second, ctx.Done())
×
1315
        go wait.Until(runWorker("update dnat for vpc nat gateway", c.updateVpcDnatQueue, c.handleUpdateVpcDnat), time.Second, ctx.Done())
×
1316
        go wait.Until(runWorker("update snat for vpc nat gateway", c.updateVpcSnatQueue, c.handleUpdateVpcSnat), time.Second, ctx.Done())
×
1317
        go wait.Until(runWorker("update subnet route for vpc nat gateway", c.updateVpcSubnetQueue, c.handleUpdateNatGwSubnetRoute), time.Second, ctx.Done())
×
1318
        go wait.Until(runWorker("add/update csr", c.addOrUpdateCsrQueue, c.handleAddOrUpdateCsr), time.Second, ctx.Done())
×
1319
        // add default and join subnet and wait them ready
×
1320
        for range c.config.WorkerNum {
×
1321
                go wait.Until(runWorker("add/update subnet", c.addOrUpdateSubnetQueue, c.handleAddOrUpdateSubnet), time.Second, ctx.Done())
×
1322
        }
×
1323
        go wait.Until(runWorker("add/update ippool", c.addOrUpdateIPPoolQueue, c.handleAddOrUpdateIPPool), time.Second, ctx.Done())
×
1324
        go wait.Until(runWorker("add vlan", c.addVlanQueue, c.handleAddVlan), time.Second, ctx.Done())
×
1325
        go wait.Until(runWorker("add namespace", c.addNamespaceQueue, c.handleAddNamespace), time.Second, ctx.Done())
×
1326
        err := wait.PollUntilContextCancel(ctx, 3*time.Second, true, func(_ context.Context) (done bool, err error) {
×
1327
                subnets := []string{c.config.DefaultLogicalSwitch, c.config.NodeSwitch}
×
1328
                klog.Infof("wait for subnets %v ready", subnets)
×
1329

×
1330
                return c.allSubnetReady(subnets...)
×
1331
        })
×
1332
        if err != nil {
×
1333
                klog.Fatalf("wait default and join subnet ready, error: %v", err)
×
1334
        }
×
1335

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

×
1340
        // run node worker before handle any pods
×
1341
        for range c.config.WorkerNum {
×
1342
                go wait.Until(runWorker("add node", c.addNodeQueue, c.handleAddNode), time.Second, ctx.Done())
×
1343
                go wait.Until(runWorker("update node", c.updateNodeQueue, c.handleUpdateNode), time.Second, ctx.Done())
×
1344
                go wait.Until(runWorker("delete node", c.deleteNodeQueue, c.handleDeleteNode), time.Second, ctx.Done())
×
1345
        }
×
1346
        for {
×
1347
                ready := true
×
1348
                time.Sleep(3 * time.Second)
×
1349
                nodes, err := c.nodesLister.List(labels.Everything())
×
1350
                if err != nil {
×
1351
                        util.LogFatalAndExit(err, "failed to list nodes")
×
1352
                }
×
1353
                for _, node := range nodes {
×
1354
                        if node.Annotations[util.AllocatedAnnotation] != "true" {
×
1355
                                klog.Infof("wait node %s annotation ready", node.Name)
×
1356
                                ready = false
×
1357
                                break
×
1358
                        }
1359
                }
1360
                if ready {
×
1361
                        break
×
1362
                }
1363
        }
1364

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

×
1370
                go wait.Until(runWorker("add/update switch lb rule", c.addSwitchLBRuleQueue, c.handleAddOrUpdateSwitchLBRule), time.Second, ctx.Done())
×
1371
                go wait.Until(runWorker("delete switch lb rule", c.delSwitchLBRuleQueue, c.handleDelSwitchLBRule), time.Second, ctx.Done())
×
1372
                go wait.Until(runWorker("delete switch lb rule", c.updateSwitchLBRuleQueue, c.handleUpdateSwitchLBRule), time.Second, ctx.Done())
×
1373

×
1374
                go wait.Until(runWorker("add/update vpc dns", c.addOrUpdateVpcDNSQueue, c.handleAddOrUpdateVPCDNS), time.Second, ctx.Done())
×
1375
                go wait.Until(runWorker("delete vpc dns", c.delVpcDNSQueue, c.handleDelVpcDNS), time.Second, ctx.Done())
×
1376
                go wait.Until(func() {
×
1377
                        c.resyncVpcDNSConfig()
×
1378
                }, 5*time.Second, ctx.Done())
×
1379
        }
1380

1381
        for range c.config.WorkerNum {
×
1382
                go wait.Until(runWorker("delete pod", c.deletePodQueue, c.handleDeletePod), time.Second, ctx.Done())
×
1383
                go wait.Until(runWorker("add/update pod", c.addOrUpdatePodQueue, c.handleAddOrUpdatePod), time.Second, ctx.Done())
×
1384
                go wait.Until(runWorker("update pod security", c.updatePodSecurityQueue, c.handleUpdatePodSecurity), time.Second, ctx.Done())
×
1385

×
1386
                go wait.Until(runWorker("delete subnet", c.deleteSubnetQueue, c.handleDeleteSubnet), time.Second, ctx.Done())
×
1387
                go wait.Until(runWorker("delete ippool", c.deleteIPPoolQueue, c.handleDeleteIPPool), time.Second, ctx.Done())
×
1388
                go wait.Until(runWorker("update status of subnet", c.updateSubnetStatusQueue, c.handleUpdateSubnetStatus), time.Second, ctx.Done())
×
1389
                go wait.Until(runWorker("update status of ippool", c.updateIPPoolStatusQueue, c.handleUpdateIPPoolStatus), time.Second, ctx.Done())
×
1390
                go wait.Until(runWorker("virtual port for subnet", c.syncVirtualPortsQueue, c.syncVirtualPort), time.Second, ctx.Done())
×
1391

×
1392
                if c.config.EnableLb {
×
1393
                        go wait.Until(runWorker("update service", c.updateServiceQueue, c.handleUpdateService), time.Second, ctx.Done())
×
1394
                        go wait.Until(runWorker("add/update endpoint slice", c.addOrUpdateEndpointSliceQueue, c.handleUpdateEndpointSlice), time.Second, ctx.Done())
×
1395
                }
×
1396

1397
                if c.config.EnableNP {
×
1398
                        go wait.Until(runWorker("update network policy", c.updateNpQueue, c.handleUpdateNp), time.Second, ctx.Done())
×
1399
                        go wait.Until(runWorker("delete network policy", c.deleteNpQueue, c.handleDeleteNp), time.Second, ctx.Done())
×
1400
                }
×
1401

1402
                go wait.Until(runWorker("delete vlan", c.delVlanQueue, c.handleDelVlan), time.Second, ctx.Done())
×
1403
                go wait.Until(runWorker("update vlan", c.updateVlanQueue, c.handleUpdateVlan), time.Second, ctx.Done())
×
1404
        }
1405

1406
        if c.config.EnableEipSnat {
×
1407
                go wait.Until(func() {
×
1408
                        // init l3 about the default vpc external lrp binding to the gw chassis
×
1409
                        c.resyncExternalGateway()
×
1410
                }, time.Second, ctx.Done())
×
1411

1412
                // maintain l3 ha about the vpc external lrp binding to the gw chassis
1413
                c.OVNNbClient.MonitorBFD()
×
1414
        }
1415
        // TODO: we should merge these two vpc nat config into one config and resync them together
1416
        go wait.Until(func() {
×
1417
                c.resyncVpcNatGwConfig()
×
1418
        }, time.Second, ctx.Done())
×
1419

1420
        go wait.Until(func() {
×
1421
                c.resyncVpcNatConfig()
×
1422
        }, time.Second, ctx.Done())
×
1423

1424
        if c.config.GCInterval != 0 {
×
1425
                go wait.Until(func() {
×
1426
                        if err := c.markAndCleanLSP(); err != nil {
×
1427
                                klog.Errorf("gc lsp error: %v", err)
×
1428
                        }
×
1429
                }, time.Duration(c.config.GCInterval)*time.Second, ctx.Done())
1430
        }
1431

1432
        go wait.Until(func() {
×
1433
                if err := c.inspectPod(); err != nil {
×
1434
                        klog.Errorf("inspection error: %v", err)
×
1435
                }
×
1436
        }, time.Duration(c.config.InspectInterval)*time.Second, ctx.Done())
1437

1438
        if c.config.EnableExternalVpc {
×
1439
                go wait.Until(func() {
×
1440
                        c.syncExternalVpc()
×
1441
                }, 5*time.Second, ctx.Done())
×
1442
        }
1443

1444
        go wait.Until(c.resyncProviderNetworkStatus, 30*time.Second, ctx.Done())
×
1445
        go wait.Until(c.exportSubnetMetrics, 30*time.Second, ctx.Done())
×
1446
        go wait.Until(c.checkSubnetGateway, 5*time.Second, ctx.Done())
×
1447
        go wait.Until(c.syncDistributedSubnetRoutes, 5*time.Second, ctx.Done())
×
1448

×
1449
        go wait.Until(runWorker("add ovn eip", c.addOvnEipQueue, c.handleAddOvnEip), time.Second, ctx.Done())
×
1450
        go wait.Until(runWorker("update ovn eip", c.updateOvnEipQueue, c.handleUpdateOvnEip), time.Second, ctx.Done())
×
1451
        go wait.Until(runWorker("reset ovn eip", c.resetOvnEipQueue, c.handleResetOvnEip), time.Second, ctx.Done())
×
1452
        go wait.Until(runWorker("delete ovn eip", c.delOvnEipQueue, c.handleDelOvnEip), time.Second, ctx.Done())
×
1453

×
1454
        go wait.Until(runWorker("add ovn fip", c.addOvnFipQueue, c.handleAddOvnFip), time.Second, ctx.Done())
×
1455
        go wait.Until(runWorker("update ovn fip", c.updateOvnFipQueue, c.handleUpdateOvnFip), time.Second, ctx.Done())
×
1456
        go wait.Until(runWorker("delete ovn fip", c.delOvnFipQueue, c.handleDelOvnFip), time.Second, ctx.Done())
×
1457

×
1458
        go wait.Until(runWorker("add ovn snat rule", c.addOvnSnatRuleQueue, c.handleAddOvnSnatRule), time.Second, ctx.Done())
×
1459
        go wait.Until(runWorker("update ovn snat rule", c.updateOvnSnatRuleQueue, c.handleUpdateOvnSnatRule), time.Second, ctx.Done())
×
1460
        go wait.Until(runWorker("delete ovn snat rule", c.delOvnSnatRuleQueue, c.handleDelOvnSnatRule), time.Second, ctx.Done())
×
1461

×
1462
        go wait.Until(runWorker("add ovn dnat", c.addOvnDnatRuleQueue, c.handleAddOvnDnatRule), time.Second, ctx.Done())
×
1463
        go wait.Until(runWorker("update ovn dnat", c.updateOvnDnatRuleQueue, c.handleUpdateOvnDnatRule), time.Second, ctx.Done())
×
1464
        go wait.Until(runWorker("delete ovn dnat", c.delOvnDnatRuleQueue, c.handleDelOvnDnatRule), time.Second, ctx.Done())
×
1465

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

×
1468
        go wait.Until(runWorker("add ip", c.addIPQueue, c.handleAddReservedIP), time.Second, ctx.Done())
×
1469
        go wait.Until(runWorker("update ip", c.updateIPQueue, c.handleUpdateIP), time.Second, ctx.Done())
×
1470
        go wait.Until(runWorker("delete ip", c.delIPQueue, c.handleDelIP), time.Second, ctx.Done())
×
1471

×
1472
        go wait.Until(runWorker("add vip", c.addVirtualIPQueue, c.handleAddVirtualIP), time.Second, ctx.Done())
×
1473
        go wait.Until(runWorker("update vip", c.updateVirtualIPQueue, c.handleUpdateVirtualIP), time.Second, ctx.Done())
×
1474
        go wait.Until(runWorker("update virtual parent for vip", c.updateVirtualParentsQueue, c.handleUpdateVirtualParents), time.Second, ctx.Done())
×
1475
        go wait.Until(runWorker("delete vip", c.delVirtualIPQueue, c.handleDelVirtualIP), time.Second, ctx.Done())
×
1476

×
1477
        go wait.Until(runWorker("add iptables eip", c.addIptablesEipQueue, c.handleAddIptablesEip), time.Second, ctx.Done())
×
1478
        go wait.Until(runWorker("update iptables eip", c.updateIptablesEipQueue, c.handleUpdateIptablesEip), time.Second, ctx.Done())
×
1479
        go wait.Until(runWorker("reset iptables eip", c.resetIptablesEipQueue, c.handleResetIptablesEip), time.Second, ctx.Done())
×
1480
        go wait.Until(runWorker("delete iptables eip", c.delIptablesEipQueue, c.handleDelIptablesEip), time.Second, ctx.Done())
×
1481

×
1482
        go wait.Until(runWorker("add iptables fip", c.addIptablesFipQueue, c.handleAddIptablesFip), time.Second, ctx.Done())
×
1483
        go wait.Until(runWorker("update iptables fip", c.updateIptablesFipQueue, c.handleUpdateIptablesFip), time.Second, ctx.Done())
×
1484
        go wait.Until(runWorker("delete iptables fip", c.delIptablesFipQueue, c.handleDelIptablesFip), time.Second, ctx.Done())
×
1485

×
1486
        go wait.Until(runWorker("add iptables dnat rule", c.addIptablesDnatRuleQueue, c.handleAddIptablesDnatRule), time.Second, ctx.Done())
×
1487
        go wait.Until(runWorker("update iptables dnat rule", c.updateIptablesDnatRuleQueue, c.handleUpdateIptablesDnatRule), time.Second, ctx.Done())
×
1488
        go wait.Until(runWorker("delete iptables dnat rule", c.delIptablesDnatRuleQueue, c.handleDelIptablesDnatRule), time.Second, ctx.Done())
×
1489

×
1490
        go wait.Until(runWorker("add iptables snat rule", c.addIptablesSnatRuleQueue, c.handleAddIptablesSnatRule), time.Second, ctx.Done())
×
1491
        go wait.Until(runWorker("update iptables snat rule", c.updateIptablesSnatRuleQueue, c.handleUpdateIptablesSnatRule), time.Second, ctx.Done())
×
1492
        go wait.Until(runWorker("delete iptables snat rule", c.delIptablesSnatRuleQueue, c.handleDelIptablesSnatRule), time.Second, ctx.Done())
×
1493

×
1494
        go wait.Until(runWorker("add qos policy", c.addQoSPolicyQueue, c.handleAddQoSPolicy), time.Second, ctx.Done())
×
1495
        go wait.Until(runWorker("update qos policy", c.updateQoSPolicyQueue, c.handleUpdateQoSPolicy), time.Second, ctx.Done())
×
1496
        go wait.Until(runWorker("delete qos policy", c.delQoSPolicyQueue, c.handleDelQoSPolicy), time.Second, ctx.Done())
×
1497

×
1498
        if c.config.EnableANP {
×
1499
                go wait.Until(runWorker("add admin network policy", c.addAnpQueue, c.handleAddAnp), time.Second, ctx.Done())
×
1500
                go wait.Until(runWorker("update admin network policy", c.updateAnpQueue, c.handleUpdateAnp), time.Second, ctx.Done())
×
1501
                go wait.Until(runWorker("delete admin network policy", c.deleteAnpQueue, c.handleDeleteAnp), time.Second, ctx.Done())
×
1502

×
1503
                go wait.Until(runWorker("add base admin network policy", c.addBanpQueue, c.handleAddBanp), time.Second, ctx.Done())
×
1504
                go wait.Until(runWorker("update base admin network policy", c.updateBanpQueue, c.handleUpdateBanp), time.Second, ctx.Done())
×
1505
                go wait.Until(runWorker("delete base admin network policy", c.deleteBanpQueue, c.handleDeleteBanp), time.Second, ctx.Done())
×
1506

×
1507
                go wait.Until(runWorker("add cluster network policy", c.addCnpQueue, c.handleAddCnp), time.Second, ctx.Done())
×
1508
                go wait.Until(runWorker("update cluster network policy", c.updateCnpQueue, c.handleUpdateCnp), time.Second, ctx.Done())
×
1509
                go wait.Until(runWorker("delete cluster network policy", c.deleteCnpQueue, c.handleDeleteCnp), time.Second, ctx.Done())
×
1510
        }
×
1511

1512
        if c.config.EnableDNSNameResolver {
×
1513
                go wait.Until(runWorker("add or update dns name resolver", c.addOrUpdateDNSNameResolverQueue, c.handleAddOrUpdateDNSNameResolver), time.Second, ctx.Done())
×
1514
                go wait.Until(runWorker("delete dns name resolver", c.deleteDNSNameResolverQueue, c.handleDeleteDNSNameResolver), time.Second, ctx.Done())
×
1515
        }
×
1516

1517
        if c.config.EnableLiveMigrationOptimize {
×
1518
                go wait.Until(runWorker("add/update vmiMigration ", c.addOrUpdateVMIMigrationQueue, c.handleAddOrUpdateVMIMigration), 50*time.Millisecond, ctx.Done())
×
1519
        }
×
1520

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

×
1523
        go wait.Until(c.dbStatus, 15*time.Second, ctx.Done())
×
1524
}
1525

1526
func (c *Controller) allSubnetReady(subnets ...string) (bool, error) {
1✔
1527
        for _, lsName := range subnets {
2✔
1528
                exist, err := c.OVNNbClient.LogicalSwitchExists(lsName)
1✔
1529
                if err != nil {
1✔
1530
                        klog.Error(err)
×
1531
                        return false, fmt.Errorf("check logical switch %s exist: %w", lsName, err)
×
1532
                }
×
1533

1534
                if !exist {
2✔
1535
                        return false, nil
1✔
1536
                }
1✔
1537
        }
1538

1539
        return true, nil
1✔
1540
}
1541

1542
func (c *Controller) initResourceOnce() {
×
1543
        c.registerSubnetMetrics()
×
1544

×
1545
        if err := c.initNodeChassis(); err != nil {
×
1546
                util.LogFatalAndExit(err, "failed to initialize node chassis")
×
1547
        }
×
1548

1549
        if err := c.initDefaultDenyAllSecurityGroup(); err != nil {
×
1550
                util.LogFatalAndExit(err, "failed to initialize 'deny_all' security group")
×
1551
        }
×
1552
        if err := c.syncSecurityGroup(); err != nil {
×
1553
                util.LogFatalAndExit(err, "failed to sync security group")
×
1554
        }
×
1555

1556
        if err := c.syncVpcNatGatewayCR(); err != nil {
×
1557
                util.LogFatalAndExit(err, "failed to sync crd vpc nat gateways")
×
1558
        }
×
1559

1560
        if err := c.initVpcNatGw(); err != nil {
×
1561
                util.LogFatalAndExit(err, "failed to initialize vpc nat gateways")
×
1562
        }
×
1563
        if c.config.EnableLb {
×
1564
                if err := c.initVpcDNSConfig(); err != nil {
×
1565
                        util.LogFatalAndExit(err, "failed to initialize vpc-dns")
×
1566
                }
×
1567
        }
1568

1569
        // remove resources in ovndb that not exist any more in kubernetes resources
1570
        // process gc at last in case of affecting other init process
1571
        if err := c.gc(); err != nil {
×
1572
                util.LogFatalAndExit(err, "failed to run gc")
×
1573
        }
×
1574
}
1575

1576
func processNextWorkItem[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error, getItemKey func(any) string) bool {
×
1577
        item, shutdown := queue.Get()
×
1578
        if shutdown {
×
1579
                return false
×
1580
        }
×
1581

1582
        err := func(item T) error {
×
1583
                defer queue.Done(item)
×
1584
                if err := handler(item); err != nil {
×
1585
                        queue.AddRateLimited(item)
×
1586
                        return fmt.Errorf("error syncing %s %q: %w, requeuing", action, getItemKey(item), err)
×
1587
                }
×
1588
                queue.Forget(item)
×
1589
                return nil
×
1590
        }(item)
1591
        if err != nil {
×
1592
                utilruntime.HandleError(err)
×
1593
                return true
×
1594
        }
×
1595
        return true
×
1596
}
1597

1598
func getWorkItemKey(obj any) string {
×
1599
        switch v := obj.(type) {
×
1600
        case string:
×
1601
                return v
×
1602
        case *vpcService:
×
1603
                return cache.MetaObjectToName(obj.(*vpcService).Svc).String()
×
1604
        case *AdminNetworkPolicyChangedDelta:
×
1605
                return v.key
×
1606
        case *SlrInfo:
×
1607
                return v.Name
×
1608
        default:
×
1609
                key, err := cache.MetaNamespaceKeyFunc(obj)
×
1610
                if err != nil {
×
1611
                        utilruntime.HandleError(err)
×
1612
                        return ""
×
1613
                }
×
1614
                return key
×
1615
        }
1616
}
1617

1618
func runWorker[T comparable](action string, queue workqueue.TypedRateLimitingInterface[T], handler func(T) error) func() {
×
1619
        return func() {
×
1620
                for processNextWorkItem(action, queue, handler, getWorkItemKey) {
×
1621
                }
×
1622
        }
1623
}
1624

1625
// apiResourceExists checks if all specified kinds exist in the given group version.
1626
// It returns true if all kinds are found, false otherwise.
1627
// Parameters:
1628
// - discoveryClient: The discovery client to use for querying API resources.
1629
// - gv: The group version string (e.g., "apps/v1").
1630
// - kinds: A variadic list of kind names to check for existence (e.g., "Deployment", "StatefulSet").
1631
func apiResourceExists(discoveryClient discovery.DiscoveryInterface, gv string, kinds ...string) (bool, error) {
×
1632
        apiResourceLists, err := discoveryClient.ServerResourcesForGroupVersion(gv)
×
1633
        if err != nil {
×
1634
                if k8serrors.IsNotFound(err) {
×
1635
                        return false, nil
×
1636
                }
×
1637
                return false, fmt.Errorf("failed to discover api resources for %s: %w", gv, err)
×
1638
        }
1639

1640
        existingKinds := set.New[string]()
×
1641
        for _, apiResource := range apiResourceLists.APIResources {
×
1642
                existingKinds.Insert(apiResource.Kind)
×
1643
        }
×
1644

1645
        return existingKinds.HasAll(kinds...), nil
×
1646
}
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