master
go 420 lines 9.95 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package k8ssd
4
5 import (
6 "context"
7 "fmt"
8 "net"
9 "strconv"
10 "strings"
11
12 "github.com/netdata/netdata/go/plugins/logger"
13 "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/model"
14
15 corev1 "k8s.io/api/core/v1"
16 "k8s.io/client-go/tools/cache"
17 "k8s.io/client-go/util/workqueue"
18 )
19
20 type podTargetGroup struct {
21 targets []model.Target
22 source string
23 }
24
25 func (p *podTargetGroup) Provider() string { return "sd:k8s:pod" }
26 func (p *podTargetGroup) Source() string { return p.source }
27 func (p *podTargetGroup) Targets() []model.Target { return p.targets }
28 func (p *podTargetGroup) setSource(src string) { p.source = src }
29
30 type PodTarget struct {
31 model.Base `hash:"ignore"`
32
33 hash uint64
34 tuid string
35
36 Address string
37 Namespace string
38 Name string
39 Annotations map[string]any
40 Labels map[string]any
41 NodeName string
42 PodIP string
43 ControllerName string
44 ControllerKind string
45 ContName string
46 Image string
47 Env map[string]any
48 Port string
49 PortName string
50 PortProtocol string
51 }
52
53 func (p PodTarget) Hash() uint64 { return p.hash }
54 func (p PodTarget) TUID() string { return p.tuid }
55
56 func newPodDiscoverer(pod, cmap, secret cache.SharedInformer) *podDiscoverer {
57
58 if pod == nil || cmap == nil || secret == nil {
59 panic("nil pod or cmap or secret informer")
60 }
61
62 queue := workqueue.NewTypedWithConfig(workqueue.TypedQueueConfig[any]{Name: "pod"})
63
64 _, _ = pod.AddEventHandler(cache.ResourceEventHandlerFuncs{
65 AddFunc: func(obj any) { enqueue(queue, obj) },
66 UpdateFunc: func(_, obj any) { enqueue(queue, obj) },
67 DeleteFunc: func(obj any) { enqueue(queue, obj) },
68 })
69
70 return &podDiscoverer{
71 Logger: log,
72 podInformer: pod,
73 cmapInformer: cmap,
74 secretInformer: secret,
75 queue: queue,
76 }
77 }
78
79 type podDiscoverer struct {
80 *logger.Logger
81 model.Base
82
83 podInformer cache.SharedInformer
84 cmapInformer cache.SharedInformer
85 secretInformer cache.SharedInformer
86 queue *workqueue.Typed[any]
87 }
88
89 func (p *podDiscoverer) String() string {
90 return "sd:k8s:pod"
91 }
92
93 func (p *podDiscoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) {
94 p.Info("instance is started")
95 defer p.Info("instance is stopped")
96 defer p.queue.ShutDown()
97
98 go p.podInformer.Run(ctx.Done())
99 go p.cmapInformer.Run(ctx.Done())
100 go p.secretInformer.Run(ctx.Done())
101
102 if !cache.WaitForCacheSync(ctx.Done(),
103 p.podInformer.HasSynced, p.cmapInformer.HasSynced, p.secretInformer.HasSynced) {
104 p.Error("failed to sync caches")
105 return
106 }
107
108 go p.run(ctx, in)
109
110 <-ctx.Done()
111 }
112
113 func (p *podDiscoverer) run(ctx context.Context, in chan<- []model.TargetGroup) {
114 for {
115 item, shutdown := p.queue.Get()
116 if shutdown {
117 return
118 }
119 p.handleQueueItem(ctx, in, item)
120 }
121 }
122
123 func (p *podDiscoverer) handleQueueItem(ctx context.Context, in chan<- []model.TargetGroup, item any) {
124 defer p.queue.Done(item)
125
126 key := item.(string)
127 namespace, name, err := cache.SplitMetaNamespaceKey(key)
128 if err != nil {
129 return
130 }
131
132 obj, ok, err := p.podInformer.GetStore().GetByKey(key)
133 if err != nil {
134 return
135 }
136
137 if !ok {
138 tgg := &podTargetGroup{source: podSourceFromNsName(namespace, name)}
139 model.SendTargetGroup(ctx, in, tgg)
140 return
141 }
142
143 pod, err := toPod(obj)
144 if err != nil {
145 return
146 }
147
148 tgg := p.buildTargetGroup(pod)
149
150 model.SendTargetGroup(ctx, in, tgg)
151
152 }
153
154 func (p *podDiscoverer) buildTargetGroup(pod *corev1.Pod) model.TargetGroup {
155 if pod.Status.PodIP == "" || len(pod.Spec.Containers) == 0 {
156 return &podTargetGroup{
157 source: podSource(pod),
158 }
159 }
160 return &podTargetGroup{
161 source: podSource(pod),
162 targets: p.buildTargets(pod),
163 }
164 }
165
166 func (p *podDiscoverer) buildTargets(pod *corev1.Pod) (targets []model.Target) {
167 var name, kind string
168 for _, ref := range pod.OwnerReferences {
169 if ref.Controller != nil && *ref.Controller {
170 name = ref.Name
171 kind = ref.Kind
172 break
173 }
174 }
175
176 for _, container := range pod.Spec.Containers {
177 env := p.collectEnv(pod.Namespace, container)
178
179 if len(container.Ports) == 0 {
180 tgt := &PodTarget{
181 tuid: podTUID(pod, container),
182 Address: pod.Status.PodIP,
183 Namespace: pod.Namespace,
184 Name: pod.Name,
185 Annotations: model.MapAny(pod.Annotations),
186 Labels: model.MapAny(pod.Labels),
187 NodeName: pod.Spec.NodeName,
188 PodIP: pod.Status.PodIP,
189 ControllerName: name,
190 ControllerKind: kind,
191 ContName: container.Name,
192 Image: container.Image,
193 Env: model.MapAny(env),
194 }
195 hash, err := model.CalcHash(tgt)
196 if err != nil {
197 continue
198 }
199 tgt.hash = hash
200
201 targets = append(targets, tgt)
202 } else {
203 for _, port := range container.Ports {
204 portNum := strconv.FormatUint(uint64(port.ContainerPort), 10)
205 tgt := &PodTarget{
206 tuid: podTUIDWithPort(pod, container, port),
207 Address: net.JoinHostPort(pod.Status.PodIP, portNum),
208 Namespace: pod.Namespace,
209 Name: pod.Name,
210 Annotations: model.MapAny(pod.Annotations),
211 Labels: model.MapAny(pod.Labels),
212 NodeName: pod.Spec.NodeName,
213 PodIP: pod.Status.PodIP,
214 ControllerName: name,
215 ControllerKind: kind,
216 ContName: container.Name,
217 Image: container.Image,
218 Env: model.MapAny(env),
219 Port: portNum,
220 PortName: port.Name,
221 PortProtocol: string(port.Protocol),
222 }
223 hash, err := model.CalcHash(tgt)
224 if err != nil {
225 continue
226 }
227 tgt.hash = hash
228
229 targets = append(targets, tgt)
230 }
231 }
232 }
233
234 return targets
235 }
236
237 func (p *podDiscoverer) collectEnv(ns string, container corev1.Container) map[string]string {
238 vars := make(map[string]string)
239
240 // When a key exists in multiple sources,
241 // the value associated with the last source will take precedence.
242 // Values defined by an Env with a duplicate key will take precedence.
243 //
244 // Order (https://github.com/kubernetes/kubectl/blob/master/pkg/describe/describe.go)
245 // - envFrom: configMapRef, secretRef
246 // - env: value || valueFrom: fieldRef, resourceFieldRef, secretRef, configMap
247
248 for _, src := range container.EnvFrom {
249 switch {
250 case src.ConfigMapRef != nil:
251 p.envFromConfigMap(vars, ns, src)
252 case src.SecretRef != nil:
253 p.envFromSecret(vars, ns, src)
254 }
255 }
256
257 for _, env := range container.Env {
258 if env.Name == "" || isVar(env.Name) {
259 continue
260 }
261 switch {
262 case env.Value != "":
263 vars[env.Name] = env.Value
264 case env.ValueFrom != nil && env.ValueFrom.SecretKeyRef != nil:
265 p.valueFromSecret(vars, ns, env)
266 case env.ValueFrom != nil && env.ValueFrom.ConfigMapKeyRef != nil:
267 p.valueFromConfigMap(vars, ns, env)
268 }
269 }
270
271 if len(vars) == 0 {
272 return nil
273 }
274 return vars
275 }
276
277 func (p *podDiscoverer) valueFromConfigMap(vars map[string]string, ns string, env corev1.EnvVar) {
278 if env.ValueFrom.ConfigMapKeyRef.Name == "" || env.ValueFrom.ConfigMapKeyRef.Key == "" {
279 return
280 }
281
282 sr := env.ValueFrom.ConfigMapKeyRef
283 key := ns + "/" + sr.Name
284
285 item, exist, err := p.cmapInformer.GetStore().GetByKey(key)
286 if err != nil || !exist {
287 return
288 }
289
290 cmap, err := toConfigMap(item)
291 if err != nil {
292 return
293 }
294
295 if v, ok := cmap.Data[sr.Key]; ok {
296 vars[env.Name] = v
297 }
298 }
299
300 func (p *podDiscoverer) valueFromSecret(vars map[string]string, ns string, env corev1.EnvVar) {
301 if env.ValueFrom.SecretKeyRef.Name == "" || env.ValueFrom.SecretKeyRef.Key == "" {
302 return
303 }
304
305 secretKey := env.ValueFrom.SecretKeyRef
306 key := ns + "/" + secretKey.Name
307
308 item, exist, err := p.secretInformer.GetStore().GetByKey(key)
309 if err != nil || !exist {
310 return
311 }
312
313 secret, err := toSecret(item)
314 if err != nil {
315 return
316 }
317
318 if v, ok := secret.Data[secretKey.Key]; ok {
319 vars[env.Name] = string(v)
320 }
321 }
322
323 func (p *podDiscoverer) envFromConfigMap(vars map[string]string, ns string, src corev1.EnvFromSource) {
324 if src.ConfigMapRef.Name == "" {
325 return
326 }
327
328 key := ns + "/" + src.ConfigMapRef.Name
329 item, exist, err := p.cmapInformer.GetStore().GetByKey(key)
330 if err != nil || !exist {
331 return
332 }
333
334 cmap, err := toConfigMap(item)
335 if err != nil {
336 return
337 }
338
339 for k, v := range cmap.Data {
340 vars[src.Prefix+k] = v
341 }
342 }
343
344 func (p *podDiscoverer) envFromSecret(vars map[string]string, ns string, src corev1.EnvFromSource) {
345 if src.SecretRef.Name == "" {
346 return
347 }
348
349 key := ns + "/" + src.SecretRef.Name
350 item, exist, err := p.secretInformer.GetStore().GetByKey(key)
351 if err != nil || !exist {
352 return
353 }
354
355 secret, err := toSecret(item)
356 if err != nil {
357 return
358 }
359
360 for k, v := range secret.Data {
361 vars[src.Prefix+k] = string(v)
362 }
363 }
364
365 func podTUID(pod *corev1.Pod, container corev1.Container) string {
366 return fmt.Sprintf("%s_%s_%s",
367 pod.Namespace,
368 pod.Name,
369 container.Name,
370 )
371 }
372
373 func podTUIDWithPort(pod *corev1.Pod, container corev1.Container, port corev1.ContainerPort) string {
374 return fmt.Sprintf("%s_%s_%s_%s_%s",
375 pod.Namespace,
376 pod.Name,
377 container.Name,
378 strings.ToLower(string(port.Protocol)),
379 strconv.FormatUint(uint64(port.ContainerPort), 10),
380 )
381 }
382
383 func podSourceFromNsName(namespace, name string) string {
384 return fmt.Sprintf("discoverer=k8s,kind=pod,namespace=%s,pod_name=%s", namespace, name)
385 }
386
387 func podSource(pod *corev1.Pod) string {
388 return podSourceFromNsName(pod.Namespace, pod.Name)
389 }
390
391 func toPod(obj any) (*corev1.Pod, error) {
392 pod, ok := obj.(*corev1.Pod)
393 if !ok {
394 return nil, fmt.Errorf("received unexpected object type: %T", obj)
395 }
396 return pod, nil
397 }
398
399 func toConfigMap(obj any) (*corev1.ConfigMap, error) {
400 cmap, ok := obj.(*corev1.ConfigMap)
401 if !ok {
402 return nil, fmt.Errorf("received unexpected object type: %T", obj)
403 }
404 return cmap, nil
405 }
406
407 func toSecret(obj any) (*corev1.Secret, error) {
408 secret, ok := obj.(*corev1.Secret)
409 if !ok {
410 return nil, fmt.Errorf("received unexpected object type: %T", obj)
411 }
412 return secret, nil
413 }
414
415 func isVar(name string) bool {
416 // Variable references $(VAR_NAME) are expanded using the previous defined
417 // environment variables in the container and any service environment
418 // variables.
419 return strings.IndexByte(name, '$') != -1
420 }