master
go 106 lines 2.27 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 newJobDiscoverer(si cache.SharedInformer, l *logger.Logger) *jobDiscoverer {
15 if si == nil {
16 panic("nil job shared informer")
17 }
18
19 queue := workqueue.NewTypedWithConfig(workqueue.TypedQueueConfig[string]{Name: "job"})
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 &jobDiscoverer{
28 Logger: l,
29 informer: si,
30 queue: queue,
31 readyCh: make(chan struct{}),
32 stopCh: make(chan struct{}),
33 }
34 }
35
36 type jobResource struct {
37 src string
38 val any
39 }
40
41 func (r jobResource) source() string { return r.src }
42 func (r jobResource) kind() kubeResourceKind { return kubeResourceJob }
43 func (r jobResource) value() any { return r.val }
44
45 type jobDiscoverer 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 *jobDiscoverer) run(ctx context.Context, in chan<- resource) {
54 d.Info("job_discoverer is started")
55 defer func() { close(d.stopCh); d.Info("job_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 *jobDiscoverer) runDiscover(ctx context.Context, in chan<- resource) {
73 for {
74 key, shutdown := d.queue.Get()
75 if shutdown {
76 return
77 }
78
79 func() {
80 defer d.queue.Done(key)
81
82 ns, name, err := cache.SplitMetaNamespaceKey(key)
83 if err != nil {
84 return
85 }
86
87 item, exists, err := d.informer.GetStore().GetByKey(key)
88 if err != nil {
89 return
90 }
91
92 r := &jobResource{src: jobSource(ns, name)}
93 if exists {
94 r.val = item
95 }
96 send(ctx, in, r)
97 }()
98 }
99 }
100
101 func jobSource(namespace, name string) string {
102 return "k8s/job/" + namespace + "/" + name
103 }
104
105 func (d *jobDiscoverer) ready() bool { return isChanClosed(d.readyCh) }
106 func (d *jobDiscoverer) stopped() bool { return isChanClosed(d.stopCh) }