master
go 207 lines 4.81 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 serviceTargetGroup struct {
21 targets []model.Target
22 source string
23 }
24
25 func (s *serviceTargetGroup) Provider() string { return "sd:k8s:service" }
26 func (s *serviceTargetGroup) Source() string { return s.source }
27 func (s *serviceTargetGroup) Targets() []model.Target { return s.targets }
28 func (s *serviceTargetGroup) setSource(src string) { s.source = src }
29
30 type ServiceTarget 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 Port string
42 PortName string
43 PortProtocol string
44 ClusterIP string
45 ExternalName string
46 Type string
47 }
48
49 func (s ServiceTarget) Hash() uint64 { return s.hash }
50 func (s ServiceTarget) TUID() string { return s.tuid }
51
52 type serviceDiscoverer struct {
53 *logger.Logger
54 model.Base
55
56 informer cache.SharedInformer
57 queue *workqueue.Typed[any]
58 }
59
60 func newServiceDiscoverer(inf cache.SharedInformer) *serviceDiscoverer {
61 if inf == nil {
62 panic("nil service informer")
63 }
64
65 queue := workqueue.NewTypedWithConfig(workqueue.TypedQueueConfig[any]{Name: "service"})
66
67 _, _ = inf.AddEventHandler(cache.ResourceEventHandlerFuncs{
68 AddFunc: func(obj any) { enqueue(queue, obj) },
69 UpdateFunc: func(_, obj any) { enqueue(queue, obj) },
70 DeleteFunc: func(obj any) { enqueue(queue, obj) },
71 })
72
73 return &serviceDiscoverer{
74 Logger: log,
75 informer: inf,
76 queue: queue,
77 }
78 }
79
80 func (s *serviceDiscoverer) String() string {
81 return "k8s service"
82 }
83
84 func (s *serviceDiscoverer) Discover(ctx context.Context, ch chan<- []model.TargetGroup) {
85 s.Info("instance is started")
86 defer s.Info("instance is stopped")
87 defer s.queue.ShutDown()
88
89 go s.informer.Run(ctx.Done())
90
91 if !cache.WaitForCacheSync(ctx.Done(), s.informer.HasSynced) {
92 s.Error("failed to sync caches")
93 return
94 }
95
96 go s.run(ctx, ch)
97
98 <-ctx.Done()
99 }
100
101 func (s *serviceDiscoverer) run(ctx context.Context, in chan<- []model.TargetGroup) {
102 for {
103 item, shutdown := s.queue.Get()
104 if shutdown {
105 return
106 }
107
108 s.handleQueueItem(ctx, in, item)
109 }
110 }
111
112 func (s *serviceDiscoverer) handleQueueItem(ctx context.Context, in chan<- []model.TargetGroup, item any) {
113 defer s.queue.Done(item)
114
115 key := item.(string)
116 namespace, name, err := cache.SplitMetaNamespaceKey(key)
117 if err != nil {
118 return
119 }
120
121 obj, exists, err := s.informer.GetStore().GetByKey(key)
122 if err != nil {
123 return
124 }
125
126 if !exists {
127 tgg := &serviceTargetGroup{source: serviceSourceFromNsName(namespace, name)}
128 model.SendTargetGroup(ctx, in, tgg)
129 return
130 }
131
132 svc, err := toService(obj)
133 if err != nil {
134 return
135 }
136
137 tgg := s.buildTargetGroup(svc)
138
139 model.SendTargetGroup(ctx, in, tgg)
140 }
141
142 func (s *serviceDiscoverer) buildTargetGroup(svc *corev1.Service) model.TargetGroup {
143 // TODO: headless service?
144 if svc.Spec.ClusterIP == "" || len(svc.Spec.Ports) == 0 {
145 return &serviceTargetGroup{
146 source: serviceSource(svc),
147 }
148 }
149 return &serviceTargetGroup{
150 source: serviceSource(svc),
151 targets: s.buildTargets(svc),
152 }
153 }
154
155 func (s *serviceDiscoverer) buildTargets(svc *corev1.Service) (targets []model.Target) {
156 for _, port := range svc.Spec.Ports {
157 portNum := strconv.FormatInt(int64(port.Port), 10)
158 tgt := &ServiceTarget{
159 tuid: serviceTUID(svc, port),
160 Address: net.JoinHostPort(svc.Name+"."+svc.Namespace+".svc", portNum),
161 Namespace: svc.Namespace,
162 Name: svc.Name,
163 Annotations: model.MapAny(svc.Annotations),
164 Labels: model.MapAny(svc.Labels),
165 Port: portNum,
166 PortName: port.Name,
167 PortProtocol: string(port.Protocol),
168 ClusterIP: svc.Spec.ClusterIP,
169 ExternalName: svc.Spec.ExternalName,
170 Type: string(svc.Spec.Type),
171 }
172 hash, err := model.CalcHash(tgt)
173 if err != nil {
174 continue
175 }
176 tgt.hash = hash
177
178 targets = append(targets, tgt)
179 }
180
181 return targets
182 }
183
184 func serviceTUID(svc *corev1.Service, port corev1.ServicePort) string {
185 return fmt.Sprintf("%s_%s_%s_%s",
186 svc.Namespace,
187 svc.Name,
188 strings.ToLower(string(port.Protocol)),
189 strconv.FormatInt(int64(port.Port), 10),
190 )
191 }
192
193 func serviceSourceFromNsName(namespace, name string) string {
194 return fmt.Sprintf("discoverer=k8s,kind=service,namespace=%s,service_name=%s", namespace, name)
195 }
196
197 func serviceSource(svc *corev1.Service) string {
198 return serviceSourceFromNsName(svc.Namespace, svc.Name)
199 }
200
201 func toService(obj any) (*corev1.Service, error) {
202 svc, ok := obj.(*corev1.Service)
203 if !ok {
204 return nil, fmt.Errorf("received unexpected object type: %T", obj)
205 }
206 return svc, nil
207 }