master
go 255 lines 6.33 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package k8ssd
4
5 import (
6 "context"
7 "fmt"
8 "log/slog"
9 "os"
10 "strings"
11 "sync"
12 "time"
13
14 "github.com/netdata/netdata/go/plugins/logger"
15 "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/model"
16 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/k8sclient"
17
18 corev1 "k8s.io/api/core/v1"
19 metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
20 "k8s.io/apimachinery/pkg/runtime"
21 "k8s.io/apimachinery/pkg/watch"
22 "k8s.io/client-go/kubernetes"
23 "k8s.io/client-go/tools/cache"
24 "k8s.io/client-go/util/workqueue"
25 )
26
27 type role string
28
29 const (
30 rolePod role = "pod"
31 roleService role = "service"
32 )
33
34 const (
35 envNodeName = "MY_NODE_NAME"
36 )
37
38 var log = logger.New().With(
39 slog.String("component", "service discovery"),
40 slog.String("discoverer", "kubernetes"),
41 )
42
43 func NewDiscoverer(cfg Config) (*KubeDiscoverer, error) {
44 if err := validateConfig(cfg); err != nil {
45 return nil, fmt.Errorf("config validation: %v", err)
46 }
47
48 client, err := k8sclient.New("Netdata/service-td")
49 if err != nil {
50 return nil, fmt.Errorf("create clientset: %v", err)
51 }
52
53 ns := cfg.Namespaces
54 if len(ns) == 0 {
55 ns = []string{corev1.NamespaceAll}
56 }
57
58 selectorField := cfg.Selector.Field
59 if role(cfg.Role) == rolePod && cfg.Pod.LocalMode {
60 name := os.Getenv(envNodeName)
61 if name == "" {
62 return nil, fmt.Errorf("local_mode is enabled, but env '%s' not set", envNodeName)
63 }
64 selectorField = joinSelectors(selectorField, "spec.nodeName="+name)
65 }
66
67 d := &KubeDiscoverer{
68 Logger: log,
69 cfgSource: cfg.Source,
70 client: client,
71 role: role(cfg.Role),
72 namespaces: ns,
73 selectorLabel: cfg.Selector.Label,
74 selectorField: selectorField,
75 discoverers: make([]model.Discoverer, 0, len(ns)),
76 started: make(chan struct{}),
77 }
78
79 return d, nil
80 }
81
82 type KubeDiscoverer struct {
83 *logger.Logger
84
85 cfgSource string
86
87 client kubernetes.Interface
88
89 role role
90 namespaces []string
91 selectorLabel string
92 selectorField string
93 discoverers []model.Discoverer
94 started chan struct{}
95 }
96
97 func (d *KubeDiscoverer) String() string {
98 return "sd:k8s"
99 }
100
101 const resyncPeriod = 10 * time.Minute
102
103 func (d *KubeDiscoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) {
104 d.Info("instance is started")
105 defer d.Info("instance is stopped")
106
107 for _, namespace := range d.namespaces {
108 var dd model.Discoverer
109 switch d.role {
110 case rolePod:
111 dd = d.setupPodDiscoverer(ctx, namespace)
112 case roleService:
113 dd = d.setupServiceDiscoverer(ctx, namespace)
114 default:
115 d.Errorf("unknown role: '%s'", d.role)
116 continue
117 }
118 d.discoverers = append(d.discoverers, dd)
119 }
120
121 if len(d.discoverers) == 0 {
122 d.Error("no discoverers registered")
123 return
124 }
125
126 d.Infof("registered: %v", d.discoverers)
127
128 var wg sync.WaitGroup
129 updates := make(chan []model.TargetGroup)
130
131 for _, disc := range d.discoverers {
132 wg.Add(1)
133 go func(disc model.Discoverer) { defer wg.Done(); disc.Discover(ctx, updates) }(disc)
134 }
135
136 done := make(chan struct{})
137 go func() { defer close(done); wg.Wait() }()
138
139 close(d.started)
140
141 for {
142 select {
143 case <-ctx.Done():
144 select {
145 case <-done:
146 d.Info("all discoverers exited")
147 case <-time.After(time.Second * 5):
148 d.Warning("not all discoverers exited")
149 }
150 return
151 case <-done:
152 d.Info("all discoverers exited")
153 return
154 case tggs := <-updates:
155 if d.cfgSource != "" {
156 for _, tgg := range tggs {
157 if v, ok := tgg.(interface{ setSource(string) }); ok {
158 src := fmt.Sprintf("%s,%s", tgg.Source(), d.cfgSource)
159 v.setSource(src)
160 }
161 }
162 }
163 select {
164 case <-ctx.Done():
165 case in <- tggs:
166 }
167 }
168 }
169 }
170
171 func (d *KubeDiscoverer) setupPodDiscoverer(ctx context.Context, ns string) *podDiscoverer {
172 pod := d.client.CoreV1().Pods(ns)
173 podLW := cache.ToListWatcherWithWatchListSemantics(&cache.ListWatch{
174 ListWithContextFunc: func(_ context.Context, opts metav1.ListOptions) (runtime.Object, error) {
175 opts.FieldSelector = d.selectorField
176 opts.LabelSelector = d.selectorLabel
177 return pod.List(ctx, opts)
178 },
179 WatchFuncWithContext: func(_ context.Context, opts metav1.ListOptions) (watch.Interface, error) {
180 opts.FieldSelector = d.selectorField
181 opts.LabelSelector = d.selectorLabel
182 return pod.Watch(ctx, opts)
183 },
184 }, d.client)
185
186 cmap := d.client.CoreV1().ConfigMaps(ns)
187 cmapLW := cache.ToListWatcherWithWatchListSemantics(&cache.ListWatch{
188 ListWithContextFunc: func(_ context.Context, opts metav1.ListOptions) (runtime.Object, error) {
189 return cmap.List(ctx, opts)
190 },
191 WatchFuncWithContext: func(_ context.Context, opts metav1.ListOptions) (watch.Interface, error) {
192 return cmap.Watch(ctx, opts)
193 },
194 }, d.client)
195
196 secret := d.client.CoreV1().Secrets(ns)
197 secretLW := cache.ToListWatcherWithWatchListSemantics(&cache.ListWatch{
198 ListWithContextFunc: func(_ context.Context, opts metav1.ListOptions) (runtime.Object, error) {
199 return secret.List(ctx, opts)
200 },
201 WatchFuncWithContext: func(_ context.Context, opts metav1.ListOptions) (watch.Interface, error) {
202 return secret.Watch(ctx, opts)
203 },
204 }, d.client)
205
206 td := newPodDiscoverer(
207 cache.NewSharedInformer(podLW, &corev1.Pod{}, resyncPeriod),
208 cache.NewSharedInformer(cmapLW, &corev1.ConfigMap{}, resyncPeriod),
209 cache.NewSharedInformer(secretLW, &corev1.Secret{}, resyncPeriod),
210 )
211
212 return td
213 }
214
215 func (d *KubeDiscoverer) setupServiceDiscoverer(ctx context.Context, namespace string) *serviceDiscoverer {
216 svc := d.client.CoreV1().Services(namespace)
217
218 svcLW := cache.ToListWatcherWithWatchListSemantics(&cache.ListWatch{
219 ListWithContextFunc: func(_ context.Context, opts metav1.ListOptions) (runtime.Object, error) {
220 opts.FieldSelector = d.selectorField
221 opts.LabelSelector = d.selectorLabel
222 return svc.List(ctx, opts)
223 },
224 WatchFuncWithContext: func(_ context.Context, opts metav1.ListOptions) (watch.Interface, error) {
225 opts.FieldSelector = d.selectorField
226 opts.LabelSelector = d.selectorLabel
227 return svc.Watch(ctx, opts)
228 },
229 }, d.client)
230
231 inf := cache.NewSharedInformer(svcLW, &corev1.Service{}, resyncPeriod)
232
233 td := newServiceDiscoverer(inf)
234
235 return td
236 }
237
238 func enqueue(queue *workqueue.Typed[any], obj any) {
239 key, err := cache.DeletionHandlingMetaNamespaceKeyFunc(obj)
240 if err != nil {
241 return
242 }
243 queue.Add(key)
244 }
245
246 func joinSelectors(srs ...string) string {
247 var i int
248 for _, v := range srs {
249 if v != "" {
250 srs[i] = v
251 i++
252 }
253 }
254 return strings.Join(srs[:i], ",")
255 }