@cryptotaxi247 / netdata-1 / commits / 197cbe5be

improvement(go.d/sd): add file path to k8s/snmp discovered job source (#19776)

Ilya Mashchenko committed Mar 5, 2025 at 13:37 UTC 197cbe5bee5b97486f1f8402bb626ef4e828e9cf
10 files changed +46 -17
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/config.go
+2
@@ -8,6 +8,8 @@ import (
8 )
9
10 type Config struct {
11 + Source string `yaml:"-"`
12 +
13 APIServer string `yaml:"api_server"` // TODO: not used
14 Role string `yaml:"role"`
15 Tags string `yaml:"tags"`
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/kubernetes.go
+11
@@ -72,6 +72,7 @@ func NewKubeDiscoverer(cfg Config) (*KubeDiscoverer, error) {
72
73 d := &KubeDiscoverer{
74 Logger: log,
75 + cfgSource: cfg.Source,
76 client: client,
77 tags: tags,
78 role: role(cfg.Role),
@@ -88,6 +89,8 @@ func NewKubeDiscoverer(cfg Config) (*KubeDiscoverer, error) {
89 type KubeDiscoverer struct {
90 *logger.Logger
91
92 + cfgSource string
93 +
94 client kubernetes.Interface
95
96 tags model.Tags
@@ -157,6 +160,14 @@ func (d *KubeDiscoverer) Discover(ctx context.Context, in chan<- []model.TargetG
160 d.Info("all discoverers exited")
161 return
162 case tggs := <-updates:
163 + if d.cfgSource != "" {
164 + for _, tgg := range tggs {
165 + if v, ok := tgg.(interface{ setSource(string) }); ok {
166 + src := fmt.Sprintf("%s,%s", tgg.Source(), d.cfgSource)
167 + v.setSource(src)
168 + }
169 + }
170 + }
171 select {
172 case <-ctx.Done():
173 case in <- tggs:
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/kubernetes_test.go
+1
@@ -137,6 +137,7 @@ func prepareDiscoverer(role role, namespaces []string, objects ...runtime.Object
137 client := fake.NewClientset(objects...)
138 tags, _ := model.ParseTags("k8s")
139 disc := &KubeDiscoverer{
140 + cfgSource: "test=test",
141 tags: tags,
142 role: role,
143 namespaces: namespaces,
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/pod.go
+4 -3
@@ -22,9 +22,10 @@ type podTargetGroup struct {
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 }
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"`
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/pod_test.go
+5 -3
@@ -44,8 +44,8 @@ func TestPodTargetGroup_Source(t *testing.T) {
44 }
45 },
46 wantSources: []string{
47 - "discoverer=k8s,kind=pod,namespace=default,pod_name=httpd-dd95c4d68-5bkwl",
48 - "discoverer=k8s,kind=pod,namespace=default,pod_name=nginx-7cfd77469b-q6kxj",
47 + "discoverer=k8s,kind=pod,namespace=default,pod_name=httpd-dd95c4d68-5bkwl,test=test",
48 + "discoverer=k8s,kind=pod,namespace=default,pod_name=nginx-7cfd77469b-q6kxj,test=test",
49 },
50 },
51 }
@@ -599,7 +599,9 @@ func prepareSecret(name string, data map[string]string) *corev1.Secret {
599 }
600
601 func prepareEmptyPodTargetGroup(pod *corev1.Pod) *podTargetGroup {
602 - return &podTargetGroup{source: podSource(pod)}
602 + tgg := &podTargetGroup{source: podSource(pod)}
603 + tgg.source += ",test=test"
604 + return tgg
605 }
606
607 func preparePodTargetGroup(pod *corev1.Pod) *podTargetGroup {
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/service.go
+4 -3
@@ -22,9 +22,10 @@ type serviceTargetGroup struct {
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 }
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"`
src/go/plugin/go.d/agent/discovery/sd/discoverer/k8ssd/service_test.go
+5 -3
@@ -43,8 +43,8 @@ func TestServiceTargetGroup_Source(t *testing.T) {
43 }
44 },
45 wantSources: []string{
46 - "discoverer=k8s,kind=service,namespace=default,service_name=httpd-cluster-ip-service",
47 - "discoverer=k8s,kind=service,namespace=default,service_name=nginx-cluster-ip-service",
46 + "discoverer=k8s,kind=service,namespace=default,service_name=httpd-cluster-ip-service,test=test",
47 + "discoverer=k8s,kind=service,namespace=default,service_name=nginx-cluster-ip-service,test=test",
48 },
49 },
50 }
@@ -425,7 +425,9 @@ func newNGINXHeadlessService() *corev1.Service {
425 }
426
427 func prepareEmptySvcTargetGroup(svc *corev1.Service) *serviceTargetGroup {
428 - return &serviceTargetGroup{source: serviceSource(svc)}
428 + tgg := &serviceTargetGroup{source: serviceSource(svc)}
429 + tgg.source += ",test=test"
430 + return tgg
431 }
432
433 func prepareSvcTargetGroup(svc *corev1.Service) *serviceTargetGroup {
src/go/plugin/go.d/agent/discovery/sd/discoverer/snmpsd/config.go
+2
@@ -13,6 +13,8 @@ import (
13
14 type (
15 Config struct {
16 + Source string `yaml:"-"`
17 +
18 // RescanInterval defines how often to scan the networks for devices (default: 30m)
19 RescanInterval *confopt.Duration `yaml:"rescan_interval"`
20 // Timeout defines the maximum time to wait for SNMP device responses (default: 1s)
src/go/plugin/go.d/agent/discovery/sd/discoverer/snmpsd/discoverer.go
+10 -5
@@ -38,9 +38,10 @@ func NewDiscoverer(cfg Config) (*Discoverer, error) {
38 slog.String("component", "service discovery"),
39 slog.String("discoverer", "snmp"),
40 ),
41 - started: make(chan struct{}),
42 - cfgHash: cfgHash,
43 - subnets: subnets,
41 + cfgSource: cfg.Source,
42 + started: make(chan struct{}),
43 + cfgHash: cfgHash,
44 + subnets: subnets,
45 newSnmpClient: func() (gosnmp.Handler, func()) {
46 return gosnmp.NewHandler(), func() {}
47 },
@@ -75,8 +76,9 @@ type (
76 *logger.Logger
77 model.Base
78
78 - started chan struct{}
79 - cfgHash uint64
79 + cfgSource string // pipeline configuration source
80 + started chan struct{}
81 + cfgHash uint64
82
83 subnets []subnet
84
@@ -164,6 +166,9 @@ func (d *Discoverer) discoverNetworks(ctx context.Context, in chan<- []model.Tar
166
167 func (d *Discoverer) discoverNetwork(ctx context.Context, in chan<- []model.TargetGroup, sub subnet, doProbing bool) {
168 tgg := newTargetGroup(sub)
169 + if d.cfgSource != "" {
170 + tgg.source += fmt.Sprintf(",%s", d.cfgSource)
171 + }
172 p := pool.New().WithMaxGoroutines(d.parallelScansPerNetwork)
173
174 for ip := range sub.ips.Iterate() {
src/go/plugin/go.d/agent/discovery/sd/pipeline/pipeline.go
+2
@@ -97,6 +97,7 @@ func (p *Pipeline) registerDiscoverers(conf Config) error {
97 p.discoverers = append(p.discoverers, td)
98 case "k8s":
99 for _, k8sCfg := range cfg.K8s {
100 + k8sCfg.Source = conf.Source
101 td, err := k8ssd.NewKubeDiscoverer(k8sCfg)
102 if err != nil {
103 return fmt.Errorf("failed to create '%s' discoverer: %v", cfg.Discoverer, err)
@@ -104,6 +105,7 @@ func (p *Pipeline) registerDiscoverers(conf Config) error {
105 p.discoverers = append(p.discoverers, td)
106 }
107 case "snmp":
108 + cfg.SNMP.Source = conf.Source
109 td, err := snmpsd.NewDiscoverer(cfg.SNMP)
110 if err != nil {
111 return fmt.Errorf("failed to create '%s' discoverer: %v", cfg.Discoverer, err)