master
go 106 lines 2.39 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package k8s_state
4
5 import (
6 "context"
7
8 "k8s.io/client-go/tools/cache"
9 "k8s.io/client-go/util/workqueue"
10
11 "github.com/netdata/netdata/go/plugins/logger"
12 )
13
14 func newDeploymentDiscoverer(si cache.SharedInformer, l *logger.Logger) *deploymentDiscoverer {
15 if si == nil {
16 panic("nil deployment& shared informer")
17 }
18
19 queue := workqueue.NewTypedWithConfig(workqueue.TypedQueueConfig[string]{Name: "replicaset"})
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 &deploymentDiscoverer{
28 Logger: l,
29 informer: si,
30 queue: queue,
31 readyCh: make(chan struct{}),
32 stopCh: make(chan struct{}),
33 }
34 }
35
36 type deployResource struct {
37 src string
38 val any
39 }
40
41 func (r deployResource) source() string { return r.src }
42 func (r deployResource) kind() kubeResourceKind { return kubeResourceDeployment }
43 func (r deployResource) value() any { return r.val }
44
45 type deploymentDiscoverer 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 *deploymentDiscoverer) run(ctx context.Context, in chan<- resource) {
54 d.Info("deployment_discoverer is started")
55 defer func() { close(d.stopCh); d.Info("deployment_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
67 close(d.readyCh)
68
69 <-ctx.Done()
70 }
71
72 func (d *deploymentDiscoverer) ready() bool { return isChanClosed(d.readyCh) }
73 func (d *deploymentDiscoverer) stopped() bool { return isChanClosed(d.stopCh) }
74
75 func (d *deploymentDiscoverer) runDiscover(ctx context.Context, in chan<- resource) {
76 for {
77 key, shutdown := d.queue.Get()
78 if shutdown {
79 return
80 }
81
82 func() {
83 defer d.queue.Done(key)
84
85 ns, name, err := cache.SplitMetaNamespaceKey(key)
86 if err != nil {
87 return
88 }
89
90 item, exists, err := d.informer.GetStore().GetByKey(key)
91 if err != nil {
92 return
93 }
94
95 r := &deployResource{src: deploymentSource(ns, name)}
96 if exists {
97 r.val = item
98 }
99 send(ctx, in, r)
100 }()
101 }
102 }
103
104 func deploymentSource(namespace, name string) string {
105 return "k8s/rs/" + namespace + "/" + name
106 }