master
go 105 lines 2.26 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package k8s_state
4
5 import (
6 "context"
7
8 "github.com/netdata/netdata/go/plugins/logger"
9
10 "k8s.io/client-go/tools/cache"
11 "k8s.io/client-go/util/workqueue"
12 )
13
14 func newNodeDiscoverer(si cache.SharedInformer, l *logger.Logger) *nodeDiscoverer {
15 if si == nil {
16 panic("nil node shared informer")
17 }
18
19 queue := workqueue.NewTypedWithConfig(workqueue.TypedQueueConfig[string]{Name: "node"})
20
21 _, _ = si.AddEventHandler(cache.ResourceEventHandlerFuncs{
22 AddFunc: func(obj any) { enqueue(queue, obj) },
23 UpdateFunc: func(_, obj any) { enqueue(queue, obj) },
24 DeleteFunc: func(obj any) { enqueue(queue, obj) },
25 })
26
27 return &nodeDiscoverer{
28 Logger: l,
29 informer: si,
30 queue: queue,
31 readyCh: make(chan struct{}),
32 stopCh: make(chan struct{}),
33 }
34 }
35
36 type nodeResource struct {
37 src string
38 val any
39 }
40
41 func (r nodeResource) source() string { return r.src }
42 func (r nodeResource) kind() kubeResourceKind { return kubeResourceNode }
43 func (r nodeResource) value() any { return r.val }
44
45 type nodeDiscoverer struct {
46 *logger.Logger
47 informer cache.SharedInformer
48 queue *workqueue.Typed[string]
49 readyCh chan struct{}
50 stopCh chan struct{}
51 }
52
53 func (d *nodeDiscoverer) run(ctx context.Context, in chan<- resource) {
54 d.Info("node_discoverer is started")
55 defer func() { close(d.stopCh); d.Info("node_discoverer is stopped") }()
56
57 defer d.queue.ShutDown()
58
59 go d.informer.Run(ctx.Done())
60
61 if !cache.WaitForCacheSync(ctx.Done(), d.informer.HasSynced) {
62 return
63 }
64
65 go d.runDiscover(ctx, in)
66 close(d.readyCh)
67
68 <-ctx.Done()
69 }
70
71 func (d *nodeDiscoverer) ready() bool { return isChanClosed(d.readyCh) }
72 func (d *nodeDiscoverer) stopped() bool { return isChanClosed(d.stopCh) }
73
74 func (d *nodeDiscoverer) runDiscover(ctx context.Context, in chan<- resource) {
75 for {
76 key, shutdown := d.queue.Get()
77 if shutdown {
78 return
79 }
80
81 func() {
82 defer d.queue.Done(key)
83
84 _, name, err := cache.SplitMetaNamespaceKey(key)
85 if err != nil {
86 return
87 }
88
89 item, exists, err := d.informer.GetStore().GetByKey(key)
90 if err != nil {
91 return
92 }
93
94 r := &nodeResource{src: nodeSource(name)}
95 if exists {
96 r.val = item
97 }
98 send(ctx, in, r)
99 }()
100 }
101 }
102
103 func nodeSource(name string) string {
104 return "k8s/node/" + name
105 }