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

kubeovn / kube-ovn / 30365630511

28 Jul 2026 01:53PM UTC coverage: 31.415%. First build
30365630511

Pull #7046

github

zhangzujian
fix: count nftables definition drift repairs

Signed-off-by: zhangzujian <zhangzujian.7@gmail.com>
Pull Request #7046: Support kube-proxy nftables mode

1244 of 1729 new or added lines in 10 files covered. (71.95%)

19858 of 63211 relevant lines covered (31.42%)

0.36 hits per line

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

17.62
/pkg/daemon/controller_linux.go
1
package daemon
2

3
import (
4
        "context"
5
        "errors"
6
        "fmt"
7
        "net"
8
        "os"
9
        "os/exec"
10
        "path/filepath"
11
        "reflect"
12
        "slices"
13
        "strings"
14
        "sync"
15
        "syscall"
16

17
        ovsutil "github.com/digitalocean/go-openvswitch/ovs"
18
        nadv1 "github.com/k8snetworkplumbingwg/network-attachment-definition-client/pkg/apis/k8s.cni.cncf.io/v1"
19
        nadutils "github.com/k8snetworkplumbingwg/network-attachment-definition-client/pkg/utils"
20
        "github.com/kubeovn/felix/ipsets"
21
        "github.com/kubeovn/go-iptables/iptables"
22
        "github.com/vishvananda/netlink"
23
        "golang.org/x/sys/unix"
24
        v1 "k8s.io/api/core/v1"
25
        k8serrors "k8s.io/apimachinery/pkg/api/errors"
26
        "k8s.io/apimachinery/pkg/labels"
27
        utilruntime "k8s.io/apimachinery/pkg/util/runtime"
28
        "k8s.io/client-go/tools/cache"
29
        "k8s.io/klog/v2"
30
        k8sipset "k8s.io/kubernetes/pkg/proxy/ipvs/ipset"
31
        k8siptables "k8s.io/kubernetes/pkg/util/iptables"
32

33
        kubeovnv1 "github.com/kubeovn/kube-ovn/pkg/apis/kubeovn/v1"
34
        "github.com/kubeovn/kube-ovn/pkg/ovs"
35
        "github.com/kubeovn/kube-ovn/pkg/util"
36
)
37

38
const (
39
        kernelModuleIPTables  = "ip_tables"
40
        kernelModuleIP6Tables = "ip6_tables"
41
)
42

43
var (
44
        setInterfaceBandwidth = ovs.SetInterfaceBandwidth
45
        configInterfaceMirror = ovs.ConfigInterfaceMirror
46
        setNetemQos           = ovs.SetNetemQos
47
)
48

49
// ControllerRuntime represents runtime specific controller members
50
type ControllerRuntime struct {
51
        iptables         map[string]*iptables.IPTables
52
        iptablesObsolete map[string]*iptables.IPTables
53
        k8siptables      map[string]k8siptables.Interface
54
        k8sipsets        k8sipset.Interface
55
        ipsets           map[string]*ipsets.IPSets
56
        gwCounters       map[string]*util.GwIPTablesCounters
57

58
        nmSyncer  *networkManagerSyncer
59
        ovsClient *ovsutil.Client
60

61
        flowCache      map[string]map[string][]string // key: bridgeName -> flowKey -> flow rules
62
        flowCacheMutex sync.RWMutex
63
        flowChan       chan struct{} // channel to trigger immediate flow sync
64
}
65

66
type LbServiceRules struct {
67
        IP          string
68
        Port        uint16
69
        Protocol    string
70
        BridgeName  string
71
        DstMac      string
72
        UnderlayNic string
73
        SubnetName  string
74
}
75

76
func evalCommandSymlinks(cmd string) (string, error) {
×
77
        path, err := exec.LookPath(cmd)
×
78
        if err != nil {
×
79
                return "", fmt.Errorf("failed to search for command %q: %w", cmd, err)
×
80
        }
×
81
        file, err := filepath.EvalSymlinks(path)
×
82
        if err != nil {
×
83
                return "", fmt.Errorf("failed to resolve symbolic links for file %q: %w", path, err)
×
84
        }
×
85

86
        return file, nil
×
87
}
88

89
func isLegacyIptablesMode() (bool, error) {
×
90
        path, err := evalCommandSymlinks("iptables")
×
91
        if err != nil {
×
92
                return false, err
×
93
        }
×
94
        pathLegacy, err := evalCommandSymlinks("iptables-legacy")
×
95
        if err != nil {
×
96
                return false, err
×
97
        }
×
98
        return path == pathLegacy, nil
×
99
}
100

101
func (c *Controller) initRuntime() error {
×
102
        ok, err := isLegacyIptablesMode()
×
103
        if err != nil {
×
104
                klog.Errorf("failed to check iptables mode: %v", err)
×
105
                return err
×
106
        }
×
107
        if !ok {
×
108
                // iptables works in nft mode, we should migrate iptables rules
×
109
                c.iptablesObsolete = make(map[string]*iptables.IPTables, 2)
×
110
        }
×
111

112
        c.iptables = make(map[string]*iptables.IPTables)
×
113
        c.ipsets = make(map[string]*ipsets.IPSets)
×
114
        c.gwCounters = make(map[string]*util.GwIPTablesCounters)
×
115
        c.k8siptables = make(map[string]k8siptables.Interface)
×
116
        c.k8sipsets = k8sipset.New()
×
117
        c.ovsClient = ovsutil.New()
×
118

×
119
        // Initialize OpenFlow flow cache (ovn-kubernetes style)
×
120
        c.flowCache = make(map[string]map[string][]string)
×
121
        c.flowChan = make(chan struct{}, 1)
×
122

×
123
        if c.protocol == kubeovnv1.ProtocolIPv4 || c.protocol == kubeovnv1.ProtocolDual {
×
124
                ipt, err := iptables.NewWithProtocol(iptables.ProtocolIPv4)
×
125
                if err != nil {
×
126
                        klog.Error(err)
×
127
                        return err
×
128
                }
×
129
                c.iptables[kubeovnv1.ProtocolIPv4] = ipt
×
130
                if c.iptablesObsolete != nil {
×
131
                        ok, err := kernelModuleLoaded(kernelModuleIPTables)
×
132
                        if err != nil {
×
133
                                klog.Errorf("failed to check kernel module %s: %v", kernelModuleIPTables, err)
×
134
                        }
×
135
                        if ok {
×
136
                                if ipt, err = iptables.NewWithProtocolAndMode(iptables.ProtocolIPv4, "legacy"); err != nil {
×
137
                                        klog.Error(err)
×
138
                                        return err
×
139
                                }
×
140
                                c.iptablesObsolete[kubeovnv1.ProtocolIPv4] = ipt
×
141
                        }
142
                }
143
                c.ipsets[kubeovnv1.ProtocolIPv4] = ipsets.NewIPSets(ipsets.NewIPVersionConfig(ipsets.IPFamilyV4, IPSetPrefix, nil, nil))
×
144
                c.k8siptables[kubeovnv1.ProtocolIPv4] = k8siptables.New(k8siptables.ProtocolIPv4)
×
145
        }
146
        if c.protocol == kubeovnv1.ProtocolIPv6 || c.protocol == kubeovnv1.ProtocolDual {
×
147
                ipt, err := iptables.NewWithProtocol(iptables.ProtocolIPv6)
×
148
                if err != nil {
×
149
                        klog.Error(err)
×
150
                        return err
×
151
                }
×
152
                c.iptables[kubeovnv1.ProtocolIPv6] = ipt
×
153
                if c.iptablesObsolete != nil {
×
154
                        ok, err := kernelModuleLoaded(kernelModuleIP6Tables)
×
155
                        if err != nil {
×
156
                                klog.Errorf("failed to check kernel module %s: %v", kernelModuleIP6Tables, err)
×
157
                        }
×
158
                        if ok {
×
159
                                if ipt, err = iptables.NewWithProtocolAndMode(iptables.ProtocolIPv6, "legacy"); err != nil {
×
160
                                        klog.Error(err)
×
161
                                        return err
×
162
                                }
×
163
                                c.iptablesObsolete[kubeovnv1.ProtocolIPv6] = ipt
×
164
                        }
165
                }
166
                c.ipsets[kubeovnv1.ProtocolIPv6] = ipsets.NewIPSets(ipsets.NewIPVersionConfig(ipsets.IPFamilyV6, IPSetPrefix, nil, nil))
×
167
                c.k8siptables[kubeovnv1.ProtocolIPv6] = k8siptables.New(k8siptables.ProtocolIPv6)
×
168
        }
169

170
        if err = ovs.ClearU2OFlows(c.ovsClient); err != nil {
×
171
                util.LogFatalAndExit(err, "failed to clear obsolete u2o flows")
×
172
        }
×
173

174
        c.nmSyncer = newNetworkManagerSyncer()
×
175
        c.nmSyncer.Run(c.transferAddrsAndRoutes)
×
176

×
177
        return nil
×
178
}
179

180
func (c *Controller) handleEnableExternalLBAddressChange(oldSubnet, newSubnet *kubeovnv1.Subnet) error {
×
181
        var subnetName string
×
182
        var action string
×
183

×
184
        switch {
×
185
        case oldSubnet != nil && newSubnet != nil:
×
186
                subnetName = oldSubnet.Name
×
187
                if oldSubnet.Spec.EnableExternalLBAddress != newSubnet.Spec.EnableExternalLBAddress {
×
188
                        klog.Infof("EnableExternalLBAddress changed for subnet %s", newSubnet.Name)
×
189
                        if newSubnet.Spec.EnableExternalLBAddress {
×
190
                                action = "add"
×
191
                        } else {
×
192
                                action = "remove"
×
193
                        }
×
194
                }
195
        case oldSubnet != nil:
×
196
                subnetName = oldSubnet.Name
×
197
                if oldSubnet.Spec.EnableExternalLBAddress {
×
198
                        klog.Infof("EnableExternalLBAddress removed for subnet %s", oldSubnet.Name)
×
199
                        action = "remove"
×
200
                }
×
201
        case newSubnet != nil:
×
202
                subnetName = newSubnet.Name
×
203
                if newSubnet.Spec.EnableExternalLBAddress {
×
204
                        klog.Infof("EnableExternalLBAddress added for subnet %s", newSubnet.Name)
×
205
                        action = "add"
×
206
                }
×
207
        }
208

209
        if action != "" {
×
210
                services, err := c.servicesLister.List(labels.Everything())
×
211
                if err != nil {
×
212
                        klog.Errorf("failed to list services: %v", err)
×
213
                        return err
×
214
                }
×
215

216
                for _, svc := range services {
×
217
                        if svc.Annotations[util.ServiceExternalIPFromSubnetAnnotation] == subnetName {
×
218
                                klog.Infof("Service %s/%s has external LB address pool annotation from subnet %s, action: %s", svc.Namespace, svc.Name, subnetName, action)
×
219
                                switch action {
×
220
                                case "add":
×
221
                                        c.serviceQueue.Add(&serviceEvent{newObj: svc})
×
222
                                case "remove":
×
223
                                        c.serviceQueue.Add(&serviceEvent{oldObj: svc})
×
224
                                }
225
                        }
226
                }
227
        }
228
        return nil
×
229
}
230

231
// handleU2OInterconnectionMACChange handles U2O interconnection MAC address changes.
232
// When U2O (Underlay to Overlay) interconnection is enabled, the svc local flow's destination
233
// MAC must point to the LRP (Logical Router Port) MAC. Otherwise, without U2O enabled (no LRP exists),
234
// the flow would hit the rules created by build_lswitch_dnat_mod_dl_dst_rules instead.
235
func (c *Controller) handleU2OInterconnectionMACChange(oldSubnet, newSubnet *kubeovnv1.Subnet) error {
×
236
        if oldSubnet == nil || newSubnet == nil {
×
237
                return nil
×
238
        }
×
239

240
        oldMAC := oldSubnet.Status.U2OInterconnectionMAC
×
241
        newMAC := newSubnet.Status.U2OInterconnectionMAC
×
242

×
243
        if oldMAC == newMAC {
×
244
                return nil
×
245
        }
×
246

247
        if newMAC == "" && oldMAC == "" {
×
248
                return nil
×
249
        }
×
250

251
        klog.Infof("U2OInterconnectionMAC changed for subnet %s: %s -> %s",
×
252
                oldSubnet.Name, oldMAC, newMAC)
×
253

×
254
        // Find all services using this subnet and re-sync them
×
255
        services, err := c.servicesLister.List(labels.Everything())
×
256
        if err != nil {
×
257
                klog.Errorf("failed to list services: %v", err)
×
258
                return err
×
259
        }
×
260

261
        for _, svc := range services {
×
262
                if svc.Annotations[util.ServiceExternalIPFromSubnetAnnotation] == oldSubnet.Name {
×
263
                        klog.Infof("Re-syncing service %s/%s due to U2OInterconnectionMAC change in subnet %s",
×
264
                                svc.Namespace, svc.Name, oldSubnet.Name)
×
265
                        c.serviceQueue.Add(&serviceEvent{newObj: svc})
×
266
                }
×
267
        }
268
        return nil
×
269
}
270

271
func (c *Controller) reconcileRouters(event *subnetEvent) error {
×
272
        subnets, err := c.subnetsLister.List(labels.Everything())
×
273
        if err != nil {
×
274
                klog.Errorf("failed to list subnets %v", err)
×
275
                return err
×
276
        }
×
277

278
        if event != nil {
×
279
                var ok bool
×
280
                var oldSubnet, newSubnet *kubeovnv1.Subnet
×
281
                if event.oldObj != nil {
×
282
                        if oldSubnet, ok = event.oldObj.(*kubeovnv1.Subnet); !ok {
×
283
                                klog.Errorf("expected old subnet in subnetEvent but got %#v", event.oldObj)
×
284
                                return nil
×
285
                        }
×
286
                }
287
                if event.newObj != nil {
×
288
                        if newSubnet, ok = event.newObj.(*kubeovnv1.Subnet); !ok {
×
289
                                klog.Errorf("expected new subnet in subnetEvent but got %#v", event.newObj)
×
290
                                return nil
×
291
                        }
×
292
                }
293

294
                if err = c.handleEnableExternalLBAddressChange(oldSubnet, newSubnet); err != nil {
×
295
                        klog.Errorf("failed to handle enable external lb address change: %v", err)
×
296
                        return err
×
297
                }
×
298

299
                if err = c.handleU2OInterconnectionMACChange(oldSubnet, newSubnet); err != nil {
×
300
                        klog.Errorf("failed to handle u2o interconnection mac change: %v", err)
×
301
                        return err
×
302
                }
×
303
                // handle policy routing
304
                rulesToAdd, rulesToDel, routesToAdd, routesToDel, err := c.diffPolicyRouting(oldSubnet, newSubnet)
×
305
                if err != nil {
×
306
                        klog.Errorf("failed to get policy routing difference: %v", err)
×
307
                        return err
×
308
                }
×
309
                // add new routes first
310
                for _, r := range routesToAdd {
×
311
                        if err = netlink.RouteReplace(&r); err != nil && !errors.Is(err, syscall.EEXIST) {
×
312
                                klog.Errorf("failed to replace route for subnet %s: %v", newSubnet.Name, err)
×
313
                                return err
×
314
                        }
×
315
                }
316
                // next, add new rules
317
                for _, r := range rulesToAdd {
×
318
                        if err = netlink.RuleAdd(&r); err != nil && !errors.Is(err, syscall.EEXIST) {
×
319
                                klog.Errorf("failed to add network rule for subnet %s: %v", newSubnet.Name, err)
×
320
                                return err
×
321
                        }
×
322
                }
323
                // then delete old network rules
324
                for _, r := range rulesToDel {
×
325
                        // loop to delete all matched rules
×
326
                        for {
×
327
                                if err = netlink.RuleDel(&r); err != nil {
×
328
                                        if !errors.Is(err, syscall.ENOENT) {
×
329
                                                klog.Errorf("failed to delete network rule for subnet %s: %v", oldSubnet.Name, err)
×
330
                                                return err
×
331
                                        }
×
332
                                        break
×
333
                                }
334
                        }
335
                }
336
                // last, delete old network routes
337
                for _, r := range routesToDel {
×
338
                        if err = netlink.RouteDel(&r); err != nil && !errors.Is(err, syscall.ENOENT) {
×
339
                                klog.Errorf("failed to delete route for subnet %s: %v", oldSubnet.Name, err)
×
340
                                return err
×
341
                        }
×
342
                }
343
        }
344

345
        node, err := c.nodesLister.Get(c.config.NodeName)
×
346
        if err != nil {
×
347
                klog.Errorf("failed to get node %s %v", c.config.NodeName, err)
×
348
                return err
×
349
        }
×
350
        nodeIPv4, nodeIPv6 := util.GetNodeInternalIP(*node)
×
351
        var joinIPv4, joinIPv6 string
×
352
        if len(node.Annotations) != 0 {
×
353
                joinIPv4, joinIPv6 = util.SplitStringIP(node.Annotations[util.IPAddressAnnotation])
×
354
        }
×
355

356
        joinCIDR := make([]string, 0, 2)
×
357
        cidrs := make([]string, 0, len(subnets)*2)
×
358
        for _, subnet := range subnets {
×
359
                // The route for overlay subnet cidr via ovn0 should not be deleted even though subnet.Status has changed to not ready
×
360
                if subnet.Spec.Vpc != c.config.ClusterRouter ||
×
361
                        (subnet.Spec.Vlan != "" && !subnet.Spec.LogicalGateway && (!subnet.Spec.U2OInterconnection || (subnet.Spec.EnableLb != nil && *subnet.Spec.EnableLb))) ||
×
362
                        !subnet.Status.IsValidated() {
×
363
                        continue
×
364
                }
365

366
                for cidrBlock := range strings.SplitSeq(subnet.Spec.CIDRBlock, ",") {
×
367
                        if _, ipNet, err := net.ParseCIDR(cidrBlock); err != nil {
×
368
                                klog.Errorf("%s is not a valid cidr block", cidrBlock)
×
369
                        } else {
×
370
                                if nodeIPv4 != "" && util.CIDRContainIP(cidrBlock, nodeIPv4) {
×
371
                                        continue
×
372
                                }
373
                                if nodeIPv6 != "" && util.CIDRContainIP(cidrBlock, nodeIPv6) {
×
374
                                        continue
×
375
                                }
376
                                cidrs = append(cidrs, ipNet.String())
×
377
                                if subnet.Name == c.config.NodeSwitch {
×
378
                                        joinCIDR = append(joinCIDR, ipNet.String())
×
379
                                }
×
380
                        }
381
                }
382
        }
383

384
        gateway, ok := node.Annotations[util.GatewayAnnotation]
×
385
        if !ok {
×
386
                err = fmt.Errorf("gateway annotation for node %s does not exist", node.Name)
×
387
                klog.Error(err)
×
388
                return err
×
389
        }
×
390
        nic, err := netlink.LinkByName(util.NodeNic)
×
391
        if err != nil {
×
392
                klog.Errorf("failed to get nic %s", util.NodeNic)
×
393
                return fmt.Errorf("failed to get nic %s", util.NodeNic)
×
394
        }
×
395

396
        allRoutes, err := getNicExistRoutes(nil, gateway)
×
397
        if err != nil {
×
398
                klog.Error(err)
×
399
                return err
×
400
        }
×
401
        nodeNicRoutes, err := getNicExistRoutes(nic, gateway)
×
402
        if err != nil {
×
403
                klog.Error(err)
×
404
                return err
×
405
        }
×
406
        toAdd, toDel := routeDiff(nodeNicRoutes, allRoutes, cidrs, joinCIDR, joinIPv4, joinIPv6, gateway, net.ParseIP(nodeIPv4), net.ParseIP(nodeIPv6))
×
407
        for _, r := range toDel {
×
408
                if err = netlink.RouteDel(&netlink.Route{Dst: r.Dst}); err != nil {
×
409
                        klog.Errorf("failed to del route %v", err)
×
410
                }
×
411
        }
412

413
        for _, r := range toAdd {
×
414
                r.LinkIndex = nic.Attrs().Index
×
415
                if err = netlink.RouteReplace(&r); err != nil {
×
416
                        klog.Errorf("failed to replace route %v: %v", r, err)
×
417
                }
×
418
        }
419

420
        return nil
×
421
}
422

423
func genLBServiceRules(service *v1.Service, bridgeName, underlayNic, dstMac, subnetName string) []LbServiceRules {
×
424
        var lbServiceRules []LbServiceRules
×
425
        for _, ingress := range service.Status.LoadBalancer.Ingress {
×
426
                for _, port := range service.Spec.Ports {
×
427
                        lbServiceRules = append(lbServiceRules, LbServiceRules{
×
428
                                IP:          ingress.IP,
×
429
                                Port:        uint16(port.Port), // #nosec G115
×
430
                                Protocol:    string(port.Protocol),
×
431
                                DstMac:      dstMac,
×
432
                                UnderlayNic: underlayNic,
×
433
                                BridgeName:  bridgeName,
×
434
                                SubnetName:  subnetName,
×
435
                        })
×
436
                }
×
437
        }
438
        return lbServiceRules
×
439
}
440

441
func (c *Controller) diffExternalLBServiceRules(oldService, newService *v1.Service, isSubnetExternalLBEnabled bool) (lbServiceRulesToAdd, lbServiceRulesToDel []LbServiceRules, err error) {
×
442
        var oldlbServiceRules, newlbServiceRules []LbServiceRules
×
443

×
444
        if oldService != nil && oldService.Annotations[util.ServiceExternalIPFromSubnetAnnotation] != "" {
×
445
                oldSubnetName := oldService.Annotations[util.ServiceExternalIPFromSubnetAnnotation]
×
446
                oldBridgeName, underlayNic, dstMac, err := c.getExtInfoBySubnet(oldSubnetName)
×
447
                if err != nil {
×
448
                        klog.Errorf("failed to get provider network by subnet %s: %v", oldSubnetName, err)
×
449
                        return nil, nil, err
×
450
                }
×
451

452
                oldlbServiceRules = genLBServiceRules(oldService, oldBridgeName, underlayNic, dstMac, oldSubnetName)
×
453
        }
454

455
        if isSubnetExternalLBEnabled && newService != nil && newService.Annotations[util.ServiceExternalIPFromSubnetAnnotation] != "" {
×
456
                newSubnetName := newService.Annotations[util.ServiceExternalIPFromSubnetAnnotation]
×
457
                newBridgeName, underlayNic, dstMac, err := c.getExtInfoBySubnet(newSubnetName)
×
458
                if err != nil {
×
459
                        klog.Errorf("failed to get provider network by subnet %s: %v", newSubnetName, err)
×
460
                        return nil, nil, err
×
461
                }
×
462
                newlbServiceRules = genLBServiceRules(newService, newBridgeName, underlayNic, dstMac, newSubnetName)
×
463
        }
464

465
        for _, oldRule := range oldlbServiceRules {
×
466
                found := slices.Contains(newlbServiceRules, oldRule)
×
467
                if !found {
×
468
                        lbServiceRulesToDel = append(lbServiceRulesToDel, oldRule)
×
469
                }
×
470
        }
471

472
        for _, newRule := range newlbServiceRules {
×
473
                found := slices.Contains(oldlbServiceRules, newRule)
×
474
                if !found {
×
475
                        lbServiceRulesToAdd = append(lbServiceRulesToAdd, newRule)
×
476
                }
×
477
        }
478

479
        return lbServiceRulesToAdd, lbServiceRulesToDel, nil
×
480
}
481

482
func (c *Controller) getExtInfoBySubnet(subnetName string) (string, string, string, error) {
×
483
        subnet, err := c.subnetsLister.Get(subnetName)
×
484
        if err != nil {
×
485
                klog.Errorf("failed to get subnet %s: %v", subnetName, err)
×
486
                return "", "", "", err
×
487
        }
×
488

489
        dstMac := subnet.Status.U2OInterconnectionMAC
×
490
        if dstMac == "" {
×
491
                dstMac = util.MasqueradeExternalLBAccessMac
×
492
                klog.Infof("Subnet %s has no U2OInterconnectionMAC, using default MAC %s", subnetName, dstMac)
×
493
        } else {
×
494
                klog.Infof("Using U2OInterconnectionMAC %s for subnet %s", dstMac, subnetName)
×
495
        }
×
496

497
        vlanName := subnet.Spec.Vlan
×
498
        if vlanName == "" {
×
499
                return "", "", "", errors.New("vlan not specified in subnet")
×
500
        }
×
501

502
        vlan, err := c.vlansLister.Get(vlanName)
×
503
        if err != nil {
×
504
                klog.Errorf("failed to get vlan %s: %v", vlanName, err)
×
505
                return "", "", "", err
×
506
        }
×
507

508
        providerNetworkName := vlan.Spec.Provider
×
509
        if providerNetworkName == "" {
×
510
                return "", "", "", errors.New("provider network not specified in vlan")
×
511
        }
×
512

513
        pn, err := c.providerNetworksLister.Get(providerNetworkName)
×
514
        if err != nil {
×
515
                klog.Errorf("failed to get provider network %s: %v", providerNetworkName, err)
×
516
                return "", "", "", err
×
517
        }
×
518

519
        underlayNic := pn.Spec.DefaultInterface
×
520
        for _, item := range pn.Spec.CustomInterfaces {
×
521
                if slices.Contains(item.Nodes, c.config.NodeName) {
×
522
                        underlayNic = item.Interface
×
523
                        break
×
524
                }
525
        }
526
        bridgeName := util.ExternalBridgeName(providerNetworkName)
×
527
        klog.Infof("Provider network: %s, Underlay NIC: %s, DstMac: %s", providerNetworkName, underlayNic, dstMac)
×
528
        return bridgeName, underlayNic, dstMac, nil
×
529
}
530

531
func (c *Controller) reconcileServices(event *serviceEvent) error {
×
532
        if event == nil {
×
533
                return nil
×
534
        }
×
535
        var ok bool
×
536
        var oldService, newService *v1.Service
×
537
        if event.oldObj != nil {
×
538
                if oldService, ok = event.oldObj.(*v1.Service); !ok {
×
539
                        klog.Errorf("expected old service in serviceEvent but got %#v", event.oldObj)
×
540
                        return nil
×
541
                }
×
542
        }
543

544
        if event.newObj != nil {
×
545
                if newService, ok = event.newObj.(*v1.Service); !ok {
×
546
                        klog.Errorf("expected new service in serviceEvent but got %#v", event.newObj)
×
547
                        return nil
×
548
                }
×
549
        }
550

551
        // check is the lb service IP related subnet's EnableExternalLBAddress
552
        isSubnetExternalLBEnabled := false
×
553
        if newService != nil && newService.Annotations[util.ServiceExternalIPFromSubnetAnnotation] != "" {
×
554
                subnet, err := c.subnetsLister.Get(newService.Annotations[util.ServiceExternalIPFromSubnetAnnotation])
×
555
                if err != nil {
×
556
                        klog.Errorf("failed to get subnet %s: %v", newService.Annotations[util.ServiceExternalIPFromSubnetAnnotation], err)
×
557
                        return err
×
558
                }
×
559
                isSubnetExternalLBEnabled = subnet.Spec.EnableExternalLBAddress
×
560
        }
561

562
        lbServiceRulesToAdd, lbServiceRulesToDel, err := c.diffExternalLBServiceRules(oldService, newService, isSubnetExternalLBEnabled)
×
563
        if err != nil {
×
564
                klog.Errorf("failed to get ip port difference: %v", err)
×
565
                return err
×
566
        }
×
567

568
        if len(lbServiceRulesToAdd) > 0 {
×
569
                for _, rule := range lbServiceRulesToAdd {
×
570
                        klog.Infof("Adding LB service rule: %+v", rule)
×
571
                        if err := c.AddOrUpdateUnderlaySubnetSvcLocalFlowCache(rule.IP, rule.Port, rule.Protocol, rule.DstMac, rule.UnderlayNic, rule.BridgeName, rule.SubnetName); err != nil {
×
572
                                klog.Errorf("failed to update underlay subnet svc local openflow cache: %v", err)
×
573
                                return err
×
574
                        }
×
575
                }
576
        }
577

578
        if len(lbServiceRulesToDel) > 0 {
×
579
                for _, rule := range lbServiceRulesToDel {
×
580
                        klog.Infof("Delete LB service rule: %+v", rule)
×
581
                        c.deleteUnderlaySubnetSvcLocalFlowCache(rule.BridgeName, rule.IP, rule.Port, rule.Protocol)
×
582
                }
×
583
        }
584

585
        return nil
×
586
}
587

588
func getNicExistRoutes(nic netlink.Link, gateway string) ([]netlink.Route, error) {
×
589
        var routes, existRoutes []netlink.Route
×
590
        var err error
×
591
        for gw := range strings.SplitSeq(gateway, ",") {
×
592
                if util.CheckProtocol(gw) == kubeovnv1.ProtocolIPv4 {
×
593
                        routes, err = netlink.RouteList(nic, netlink.FAMILY_V4)
×
594
                } else {
×
595
                        routes, err = netlink.RouteList(nic, netlink.FAMILY_V6)
×
596
                }
×
597
                if err != nil {
×
598
                        return nil, err
×
599
                }
×
600
                existRoutes = append(existRoutes, routes...)
×
601
        }
602
        return existRoutes, nil
×
603
}
604

605
func routeDiff(nodeNicRoutes, allRoutes []netlink.Route, cidrs, joinCIDR []string, joinIPv4, joinIPv6, gateway string, srcIPv4, srcIPv6 net.IP) (toAdd, toDel []netlink.Route) {
×
606
        // joinIPv6 is not used for now
×
607
        _ = joinIPv6
×
608

×
609
        for _, route := range nodeNicRoutes {
×
610
                if route.Scope == netlink.SCOPE_LINK || route.Dst == nil || route.Dst.IP.IsLinkLocalUnicast() {
×
611
                        continue
×
612
                }
613

614
                found := slices.Contains(cidrs, route.Dst.String())
×
615
                if !found {
×
616
                        toDel = append(toDel, route)
×
617
                }
×
618
                conflict := false
×
619
                for _, ar := range allRoutes {
×
620
                        if ar.Dst != nil && ar.Dst.String() == route.Dst.String() && ar.LinkIndex != route.LinkIndex {
×
621
                                // route conflict
×
622
                                conflict = true
×
623
                                break
×
624
                        }
625
                }
626
                if conflict {
×
627
                        toDel = append(toDel, route)
×
628
                }
×
629
        }
630
        if len(toDel) > 0 {
×
631
                klog.Infof("routes to delete: %v", toDel)
×
632
        }
×
633

634
        ipv4, ipv6 := util.SplitStringIP(gateway)
×
635
        gwV4, gwV6 := net.ParseIP(ipv4), net.ParseIP(ipv6)
×
636
        for _, c := range cidrs {
×
637
                var src, gw net.IP
×
638
                switch util.CheckProtocol(c) {
×
639
                case kubeovnv1.ProtocolIPv4:
×
640
                        src, gw = srcIPv4, gwV4
×
641
                case kubeovnv1.ProtocolIPv6:
×
642
                        src, gw = srcIPv6, gwV6
×
643
                }
644

645
                found := false
×
646
                for _, ar := range allRoutes {
×
647
                        if ar.Dst != nil && ar.Dst.String() == c {
×
648
                                if slices.Contains(joinCIDR, c) {
×
649
                                        // Only compare Dst for join subnets
×
650
                                        found = true
×
651
                                        klog.V(3).Infof("[routeDiff] joinCIDR route already exists in allRoutes: %v", ar)
×
652
                                        break
×
653
                                } else if (ar.Src == nil && src == nil) || (ar.Src != nil && src != nil && ar.Src.Equal(src)) {
×
654
                                        // For non-join subnets, both Dst and Src must be the same
×
655
                                        found = true
×
656
                                        klog.V(3).Infof("[routeDiff] route already exists in allRoutes: %v", ar)
×
657
                                        break
×
658
                                }
659
                        }
660
                }
661
                if found {
×
662
                        continue
×
663
                }
664
                for _, r := range nodeNicRoutes {
×
665
                        if r.Dst == nil || r.Dst.String() != c {
×
666
                                continue
×
667
                        }
668
                        if (src == nil && r.Src == nil) || (src != nil && r.Src != nil && src.Equal(r.Src)) {
×
669
                                found = true
×
670
                                break
×
671
                        }
672
                }
673
                if !found {
×
674
                        var priority int
×
675
                        scope := netlink.SCOPE_UNIVERSE
×
676
                        proto := netlink.RouteProtocol(syscall.RTPROT_STATIC)
×
677
                        if slices.Contains(joinCIDR, c) {
×
678
                                if util.CheckProtocol(c) == kubeovnv1.ProtocolIPv4 {
×
679
                                        src = net.ParseIP(joinIPv4)
×
680
                                } else {
×
681
                                        src, priority = nil, 256
×
682
                                }
×
683
                                gw, scope = nil, netlink.SCOPE_LINK
×
684
                                proto = netlink.RouteProtocol(unix.RTPROT_KERNEL)
×
685
                        }
686
                        _, cidr, _ := net.ParseCIDR(c)
×
687
                        toAdd = append(toAdd, netlink.Route{
×
688
                                Dst:      cidr,
×
689
                                Src:      src,
×
690
                                Gw:       gw,
×
691
                                Protocol: proto,
×
692
                                Scope:    scope,
×
693
                                Priority: priority,
×
694
                        })
×
695
                }
696
        }
697
        if len(toAdd) > 0 {
×
698
                klog.Infof("routes to add: %v", toAdd)
×
699
        }
×
700
        return toAdd, toDel
×
701
}
702

703
func getRulesToAdd(oldRules, newRules []netlink.Rule) []netlink.Rule {
×
704
        var toAdd []netlink.Rule
×
705

×
706
        for _, rule := range newRules {
×
707
                var found bool
×
708
                for _, r := range oldRules {
×
709
                        if r.Family == rule.Family && r.Priority == rule.Priority && r.Table == rule.Table && reflect.DeepEqual(r.Src, rule.Src) {
×
710
                                found = true
×
711
                                break
×
712
                        }
713
                }
714
                if !found {
×
715
                        toAdd = append(toAdd, rule)
×
716
                }
×
717
        }
718

719
        return toAdd
×
720
}
721

722
func getRoutesToAdd(oldRoutes, newRoutes []netlink.Route) []netlink.Route {
×
723
        var toAdd []netlink.Route
×
724

×
725
        for _, route := range newRoutes {
×
726
                var found bool
×
727
                for _, r := range oldRoutes {
×
728
                        if r.Equal(route) {
×
729
                                found = true
×
730
                                break
×
731
                        }
732
                }
733
                if !found {
×
734
                        toAdd = append(toAdd, route)
×
735
                }
×
736
        }
737

738
        return toAdd
×
739
}
740

741
func (c *Controller) diffPolicyRouting(oldSubnet, newSubnet *kubeovnv1.Subnet) (rulesToAdd, rulesToDel []netlink.Rule, routesToAdd, routesToDel []netlink.Route, err error) {
×
742
        oldRules, oldRoutes, err := c.getPolicyRouting(oldSubnet)
×
743
        if err != nil {
×
744
                klog.Error(err)
×
745
                return rulesToAdd, rulesToDel, routesToAdd, routesToDel, err
×
746
        }
×
747
        newRules, newRoutes, err := c.getPolicyRouting(newSubnet)
×
748
        if err != nil {
×
749
                klog.Error(err)
×
750
                return rulesToAdd, rulesToDel, routesToAdd, routesToDel, err
×
751
        }
×
752

753
        rulesToAdd = getRulesToAdd(oldRules, newRules)
×
754
        rulesToDel = getRulesToAdd(newRules, oldRules)
×
755
        routesToAdd = getRoutesToAdd(oldRoutes, newRoutes)
×
756
        routesToDel = getRoutesToAdd(newRoutes, oldRoutes)
×
757

×
758
        return rulesToAdd, rulesToDel, routesToAdd, routesToDel, err
×
759
}
760

761
func (c *Controller) getPolicyRouting(subnet *kubeovnv1.Subnet) ([]netlink.Rule, []netlink.Route, error) {
1✔
762
        if subnet == nil || subnet.Spec.ExternalEgressGateway == "" || subnet.Spec.Vpc != c.config.ClusterRouter {
2✔
763
                return nil, nil, nil
1✔
764
        }
1✔
765
        if subnet.Spec.GatewayType == kubeovnv1.GWCentralizedType {
2✔
766
                node, err := c.nodesLister.Get(c.config.NodeName)
1✔
767
                if err != nil {
1✔
768
                        klog.Errorf("failed to get node %s: %v", c.config.NodeName, err)
×
769
                        return nil, nil, err
×
770
                }
×
771
                isGatewayNode := util.GatewayContains(subnet.Spec.GatewayNode, c.config.NodeName) ||
1✔
772
                        (subnet.Spec.GatewayNode == "" && util.MatchLabelSelectors(subnet.Spec.GatewayNodeSelectors, node.Labels))
1✔
773
                if !isGatewayNode {
1✔
774
                        return nil, nil, nil
×
775
                }
×
776
        }
777

778
        protocols := make([]string, 1, 2)
1✔
779
        if protocol := util.CheckProtocol(subnet.Spec.ExternalEgressGateway); protocol == kubeovnv1.ProtocolDual {
2✔
780
                protocols[0] = kubeovnv1.ProtocolIPv4
1✔
781
                protocols = append(protocols, kubeovnv1.ProtocolIPv6)
1✔
782
        } else {
2✔
783
                protocols[0] = protocol
1✔
784
        }
1✔
785

786
        cidr := strings.Split(subnet.Spec.CIDRBlock, ",")
1✔
787
        egw := util.SplitTrimmed(subnet.Spec.ExternalEgressGateway, ",")
1✔
788
        if len(egw) == 0 {
1✔
789
                return nil, nil, nil
×
790
        }
×
791

792
        // rules
793
        var rules []netlink.Rule
1✔
794
        rule := netlink.NewRule()
1✔
795
        rule.Table = int(subnet.Spec.PolicyRoutingTableID)
1✔
796
        rule.Priority = int(subnet.Spec.PolicyRoutingPriority)
1✔
797
        if subnet.Spec.GatewayType == kubeovnv1.GWDistributedType {
2✔
798
                pods, err := c.podsLister.List(labels.Everything())
1✔
799
                if err != nil {
1✔
800
                        klog.Errorf("list pods failed, %+v", err)
×
801
                        return nil, nil, err
×
802
                }
×
803

804
                for _, pod := range pods {
2✔
805
                        if pod.Status.PodIP == "" ||
1✔
806
                                pod.Annotations[util.LogicalSwitchAnnotation] != subnet.Name {
2✔
807
                                continue
1✔
808
                        }
809

810
                        for i := range protocols {
2✔
811
                                rule.Family, _ = util.ProtocolToFamily(protocols[i])
1✔
812

1✔
813
                                var ip net.IP
1✔
814
                                var maskBits int
1✔
815
                                if len(pod.Status.PodIPs) == 2 && protocols[i] == kubeovnv1.ProtocolIPv6 {
2✔
816
                                        ip = net.ParseIP(pod.Status.PodIPs[1].IP)
1✔
817
                                        maskBits = 128
1✔
818
                                } else if util.CheckProtocol(pod.Status.PodIP) == protocols[i] {
3✔
819
                                        ip = net.ParseIP(pod.Status.PodIP)
1✔
820
                                        maskBits = 32
1✔
821
                                        if rule.Family == netlink.FAMILY_V6 {
2✔
822
                                                maskBits = 128
1✔
823
                                        }
1✔
824
                                }
825

826
                                if ip == nil {
2✔
827
                                        continue
1✔
828
                                }
829
                                rule.Src = &net.IPNet{IP: ip, Mask: net.CIDRMask(maskBits, maskBits)}
1✔
830
                                rules = append(rules, *rule)
1✔
831
                        }
832
                }
833
        } else {
1✔
834
                for i := range protocols {
2✔
835
                        rule.Family, _ = util.ProtocolToFamily(protocols[i])
1✔
836
                        if i >= len(cidr) {
2✔
837
                                continue
1✔
838
                        }
839
                        _, ipNet, err := net.ParseCIDR(cidr[i])
1✔
840
                        if err != nil {
1✔
841
                                klog.Errorf("failed to parse CIDR %q for subnet %s policy routing: %v", cidr[i], subnet.Name, err)
×
842
                                continue
×
843
                        }
844
                        rule.Src = ipNet
1✔
845
                        rules = append(rules, *rule)
1✔
846
                }
847
        }
848

849
        // routes
850
        var routes []netlink.Route
1✔
851
        for i := range protocols {
2✔
852
                routes = append(routes, netlink.Route{
1✔
853
                        Protocol: netlink.RouteProtocol(syscall.RTPROT_STATIC),
1✔
854
                        Table:    int(subnet.Spec.PolicyRoutingTableID),
1✔
855
                        Gw:       net.ParseIP(egw[i]),
1✔
856
                })
1✔
857
        }
1✔
858

859
        return rules, routes, nil
1✔
860
}
861

862
func (c *Controller) handleUpdatePod(key string) error {
1✔
863
        namespace, name, err := cache.SplitMetaNamespaceKey(key)
1✔
864
        if err != nil {
1✔
865
                utilruntime.HandleError(fmt.Errorf("invalid resource key: %s", key))
×
866
                return nil
×
867
        }
×
868
        klog.Infof("handle qos update for pod %s/%s", namespace, name)
1✔
869

1✔
870
        pod, err := c.podsLister.Pods(namespace).Get(name)
1✔
871
        if err != nil {
2✔
872
                if k8serrors.IsNotFound(err) {
2✔
873
                        return nil
1✔
874
                }
1✔
875
                klog.Error(err)
×
876
                return err
×
877
        }
878

879
        if err := util.ValidatePodNetwork(pod.Annotations); err != nil {
2✔
880
                klog.Errorf("validate pod %s/%s failed, %v", namespace, name, err)
1✔
881
                c.recorder.Eventf(pod, v1.EventTypeWarning, "ValidatePodNetworkFailed", "%s", err.Error())
1✔
882
                return err
1✔
883
        }
1✔
884

885
        podName := pod.Name
1✔
886
        if pod.Annotations[fmt.Sprintf(util.VMAnnotationTemplate, util.OvnProvider)] != "" {
1✔
887
                podName = pod.Annotations[fmt.Sprintf(util.VMAnnotationTemplate, util.OvnProvider)]
×
888
        }
×
889

890
        // set default nic bandwidth
891
        //  ovsIngress and ovsEgress are derived from the pod's egress and ingress rate annotations respectively, their roles are reversed from the OVS interface perspective.
892
        ifaceID := ovs.PodNameToPortName(podName, pod.Namespace, util.OvnProvider)
1✔
893
        ovsIngress := pod.Annotations[util.EgressRateAnnotation]
1✔
894
        ovsEgress := pod.Annotations[util.IngressRateAnnotation]
1✔
895
        ovsIngressBurst := pod.Annotations[util.EgressBurstAnnotation]
1✔
896
        ovsEgressBurst := pod.Annotations[util.IngressBurstAnnotation]
1✔
897
        err = setInterfaceBandwidth(podName, pod.Namespace, ifaceID, ovsIngress, ovsEgress, ovsIngressBurst, ovsEgressBurst)
1✔
898
        if err != nil {
2✔
899
                klog.Error(err)
1✔
900
                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=bandwidth provider=%s interface=%s node=%s: %v", util.OvnProvider, ifaceID, c.config.NodeName, err)
1✔
901
                return err
1✔
902
        }
1✔
903
        err = configInterfaceMirror(c.config.EnableMirror, pod.Annotations[util.MirrorControlAnnotation], ifaceID)
1✔
904
        if err != nil {
1✔
905
                klog.Error(err)
×
906
                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=mirror provider=%s interface=%s node=%s: %v", util.OvnProvider, ifaceID, c.config.NodeName, err)
×
907
                return err
×
908
        }
×
909
        // set linux-netem qos
910
        err = setNetemQos(podName, pod.Namespace, ifaceID, pod.Annotations[util.NetemQosLatencyAnnotation], pod.Annotations[util.NetemQosLimitAnnotation], pod.Annotations[util.NetemQosLossAnnotation], pod.Annotations[util.NetemQosJitterAnnotation])
1✔
911
        if err != nil {
1✔
912
                klog.Error(err)
×
913
                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=netem provider=%s interface=%s node=%s: %v", util.OvnProvider, ifaceID, c.config.NodeName, err)
×
914
                return err
×
915
        }
×
916
        processed := []string{fmt.Sprintf("provider=%s interface=%s", util.OvnProvider, ifaceID)}
1✔
917

1✔
918
        // set multus-nic bandwidth
1✔
919
        attachNets, err := nadutils.ParsePodNetworkAnnotation(pod)
1✔
920
        if err != nil {
2✔
921
                if _, ok := err.(*nadv1.NoK8sNetworkError); !ok {
2✔
922
                        klog.Error(err)
1✔
923
                        c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=parseNetworkAttachment provider=unknown interface=unknown node=%s: %v", c.config.NodeName, err)
1✔
924
                        return err
1✔
925
                }
1✔
926
        }
927
        for _, multiNet := range attachNets {
2✔
928
                provider := fmt.Sprintf("%s.%s.%s", multiNet.Name, multiNet.Namespace, util.OvnProvider)
1✔
929
                multiNetPodName := pod.Name
1✔
930
                if pod.Annotations[fmt.Sprintf(util.VMAnnotationTemplate, provider)] != "" {
2✔
931
                        multiNetPodName = pod.Annotations[fmt.Sprintf(util.VMAnnotationTemplate, provider)]
1✔
932
                }
1✔
933
                if pod.Annotations[fmt.Sprintf(util.AllocatedAnnotationTemplate, provider)] == "true" {
2✔
934
                        ifaceID = ovs.PodNameToPortName(multiNetPodName, pod.Namespace, provider)
1✔
935
                        err = setInterfaceBandwidth(multiNetPodName, pod.Namespace, ifaceID,
1✔
936
                                pod.Annotations[fmt.Sprintf(util.EgressRateAnnotationTemplate, provider)],
1✔
937
                                pod.Annotations[fmt.Sprintf(util.IngressRateAnnotationTemplate, provider)],
1✔
938
                                pod.Annotations[fmt.Sprintf(util.EgressBurstAnnotationTemplate, provider)],
1✔
939
                                pod.Annotations[fmt.Sprintf(util.IngressBurstAnnotationTemplate, provider)])
1✔
940
                        if err != nil {
2✔
941
                                klog.Error(err)
1✔
942
                                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=bandwidth provider=%s interface=%s node=%s: %v", provider, ifaceID, c.config.NodeName, err)
1✔
943
                                return err
1✔
944
                        }
1✔
945
                        err = configInterfaceMirror(c.config.EnableMirror, pod.Annotations[fmt.Sprintf(util.MirrorControlAnnotationTemplate, provider)], ifaceID)
1✔
946
                        if err != nil {
2✔
947
                                klog.Error(err)
1✔
948
                                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=mirror provider=%s interface=%s node=%s: %v", provider, ifaceID, c.config.NodeName, err)
1✔
949
                                return err
1✔
950
                        }
1✔
951
                        err = setNetemQos(multiNetPodName, pod.Namespace, ifaceID, pod.Annotations[fmt.Sprintf(util.NetemQosLatencyAnnotationTemplate, provider)], pod.Annotations[fmt.Sprintf(util.NetemQosLimitAnnotationTemplate, provider)], pod.Annotations[fmt.Sprintf(util.NetemQosLossAnnotationTemplate, provider)], pod.Annotations[fmt.Sprintf(util.NetemQosJitterAnnotationTemplate, provider)])
1✔
952
                        if err != nil {
2✔
953
                                klog.Error(err)
1✔
954
                                c.recorder.Eventf(pod, v1.EventTypeWarning, "PodQoSUpdateFailed", "Failed to update pod QoS: stage=netem provider=%s interface=%s node=%s: %v", provider, ifaceID, c.config.NodeName, err)
1✔
955
                                return err
1✔
956
                        }
1✔
957
                        processed = append(processed, fmt.Sprintf("provider=%s interface=%s", provider, ifaceID))
1✔
958
                }
959
        }
960

961
        c.recorder.Eventf(pod, v1.EventTypeNormal, "PodQoSUpdated", "Updated pod QoS: node=%s %s", c.config.NodeName, strings.Join(processed, ", "))
1✔
962
        return nil
1✔
963
}
964

965
func (c *Controller) loopEncapIPCheck() {
×
966
        node, err := c.nodesLister.Get(c.config.NodeName)
×
967
        if err != nil {
×
968
                klog.Errorf("failed to get node %s %v", c.config.NodeName, err)
×
969
                return
×
970
        }
×
971

972
        if nodeTunnelName := node.GetAnnotations()[util.TunnelInterfaceAnnotation]; nodeTunnelName != "" {
×
973
                iface, err := findInterface(nodeTunnelName)
×
974
                if err != nil {
×
975
                        klog.Errorf("failed to find iface %s, %v", nodeTunnelName, err)
×
976
                        return
×
977
                }
×
978
                if iface.Flags&net.FlagUp == 0 {
×
979
                        klog.Errorf("iface %v is down", nodeTunnelName)
×
980
                        return
×
981
                }
×
982
                addrs, err := iface.Addrs()
×
983
                if err != nil {
×
984
                        klog.Errorf("failed to get iface addr. %v", err)
×
985
                        return
×
986
                }
×
987
                if len(addrs) == 0 {
×
988
                        klog.Errorf("iface %s has no ip address", nodeTunnelName)
×
989
                        return
×
990
                }
×
991
                if iface.Name != c.config.tunnelIface {
×
992
                        klog.Infof("use %s as tunnel interface", iface.Name)
×
993
                        c.config.tunnelIface = iface.Name
×
994
                }
×
995

996
                // if assigned iface in node annotation is down or with no ip, the error msg should be printed periodically
997
                if c.config.Iface == nodeTunnelName {
×
998
                        klog.V(3).Infof("node tunnel interface %s not changed", nodeTunnelName)
×
999
                        return
×
1000
                }
×
1001

1002
                var encapIP string
×
1003
                for _, addr := range addrs {
×
1004
                        ipStr, _, _ := strings.Cut(addr.String(), "/")
×
1005
                        if ip := net.ParseIP(ipStr); ip == nil || ip.IsLinkLocalUnicast() || ip.IsLoopback() {
×
1006
                                continue
×
1007
                        }
1008
                        encapIP = ipStr
×
1009
                        break
×
1010
                }
1011
                if encapIP == "" {
×
1012
                        klog.Errorf("iface %s has no valid IP address", nodeTunnelName)
×
1013
                        return
×
1014
                }
×
1015

1016
                c.config.Iface = nodeTunnelName
×
1017
                klog.Infof("Update node tunnel interface %v", nodeTunnelName)
×
1018
                c.config.SetDefaultEncapIP(encapIP)
×
1019
                if err = c.config.setEncapIPs(); err != nil {
×
1020
                        klog.Errorf("failed to set encap ip %s for iface %s", encapIP, nodeTunnelName)
×
1021
                        return
×
1022
                }
×
1023
        }
1024
}
1025

1026
func (c *Controller) ovnMetricsUpdate() {
×
NEW
1027
        if err := c.gatewayBackendManager.ReadSubnetCounters(context.Background()); err != nil {
×
NEW
1028
                klog.Errorf("failed to read gateway subnet counters: %v", err)
×
NEW
1029
        }
×
1030

1031
        resetSysParaMetrics()
×
1032
        c.setIPLocalPortRangeMetric()
×
1033
        c.setCheckSumErrMetric()
×
1034
        c.setDNSSearchMetric()
×
1035
        c.setTCPTwRecycleMetric()
×
1036
        c.setTCPMtuProbingMetric()
×
1037
        c.setConntrackTCPLiberalMetric()
×
1038
        c.setBridgeNfCallIptablesMetric()
×
1039
        c.setIPv6RouteMaxsizeMetric()
×
1040
        c.setTCPMemMetric()
×
1041
}
1042

1043
func resetSysParaMetrics() {
×
1044
        metricIPLocalPortRange.Reset()
×
1045
        metricCheckSumErr.Reset()
×
1046
        metricDNSSearch.Reset()
×
1047
        metricTCPTwRecycle.Reset()
×
1048
        metricTCPMtuProbing.Reset()
×
1049
        metricConntrackTCPLiberal.Reset()
×
1050
        metricBridgeNfCallIptables.Reset()
×
1051
        metricTCPMem.Reset()
×
1052
        metricIPv6RouteMaxsize.Reset()
×
1053
}
×
1054

1055
func rotateLog() {
×
1056
        output, err := exec.Command("logrotate", "/etc/logrotate.d/openvswitch").CombinedOutput()
×
1057
        if err != nil {
×
1058
                klog.Errorf("failed to rotate openvswitch log %q", output)
×
1059
        }
×
1060
        output, err = exec.Command("logrotate", "/etc/logrotate.d/ovn").CombinedOutput()
×
1061
        if err != nil {
×
1062
                klog.Errorf("failed to rotate ovn log %q", output)
×
1063
        }
×
1064
        output, err = exec.Command("logrotate", "/etc/logrotate.d/kubeovn").CombinedOutput()
×
1065
        if err != nil {
×
1066
                klog.Errorf("failed to rotate kube-ovn log %q", output)
×
1067
        }
×
1068
}
1069

1070
func kernelModuleLoaded(module string) (bool, error) {
×
1071
        data, err := os.ReadFile("/proc/modules")
×
1072
        if err != nil {
×
1073
                klog.Errorf("failed to read /proc/modules: %v", err)
×
1074
                return false, err
×
1075
        }
×
1076

1077
        for line := range strings.SplitSeq(string(data), "\n") {
×
1078
                if fields := strings.Fields(line); len(fields) != 0 && fields[0] == module {
×
1079
                        return true, nil
×
1080
                }
×
1081
        }
1082

1083
        return false, nil
×
1084
}
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