master
go 2,357 lines 74.5 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package sd
4
5 import (
6 "bytes"
7 "context"
8 "encoding/json"
9 "errors"
10 "fmt"
11 "strings"
12 "testing"
13 "time"
14
15 "github.com/netdata/netdata/go/plugins/logger"
16 "github.com/netdata/netdata/go/plugins/pkg/confopt"
17 "github.com/netdata/netdata/go/plugins/pkg/netdataapi"
18 "github.com/netdata/netdata/go/plugins/pkg/safewriter"
19 "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/pipeline"
20 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
21 "github.com/netdata/netdata/go/plugins/plugin/framework/dyncfg"
22 "github.com/netdata/netdata/go/plugins/plugin/framework/functions"
23
24 "github.com/stretchr/testify/assert"
25 "github.com/stretchr/testify/require"
26 )
27
28 // Helper functions to create test configs using pipeline.Config
29
30 // defaultTestServices returns a minimal valid service rule for tests.
31 func defaultTestServices() []pipeline.ServiceRuleConfig {
32 return []pipeline.ServiceRuleConfig{
33 {ID: "test-rule", Match: "true"},
34 }
35 }
36
37 func mustDiscovererPayload(typ string, cfg any) pipeline.DiscovererPayload {
38 p, err := pipeline.NewDiscovererPayload(typ, cfg)
39 if err != nil {
40 panic(err)
41 }
42 return p
43 }
44
45 func newTestNetListenersConfig(name string, interval confopt.LongDuration, timeout confopt.Duration, services []pipeline.ServiceRuleConfig) pipeline.Config {
46 return pipeline.Config{
47 Name: name,
48 Discoverer: mustDiscovererPayload(testDiscovererTypeNetListeners, testNetListenersConfig{
49 Interval: interval,
50 Timeout: timeout,
51 }),
52 Services: services,
53 }
54 }
55
56 func newTestDockerConfig(name, address string, timeout confopt.Duration, services []pipeline.ServiceRuleConfig) pipeline.Config {
57 return pipeline.Config{
58 Name: name,
59 Discoverer: mustDiscovererPayload(testDiscovererTypeDocker, testDockerConfig{
60 Address: address,
61 Timeout: timeout,
62 }),
63 Services: services,
64 }
65 }
66
67 func newTestK8sConfig(name string, cfgs []testK8sConfig, services []pipeline.ServiceRuleConfig) pipeline.Config {
68 return pipeline.Config{
69 Name: name,
70 Discoverer: mustDiscovererPayload(testDiscovererTypeK8s, cfgs),
71 Services: services,
72 }
73 }
74
75 func newTestSNMPConfig(name string, cfg testSNMPConfig, services []pipeline.ServiceRuleConfig) pipeline.Config {
76 return pipeline.Config{
77 Name: name,
78 Discoverer: mustDiscovererPayload(testDiscovererTypeSNMP, cfg),
79 Services: services,
80 }
81 }
82
83 type dyncfgSim struct {
84 do func(sd *ServiceDiscovery)
85
86 wantExposed []wantExposedConfig
87 wantRunning []string
88 wantDyncfg string
89 wantDyncfgFunc func(t *testing.T, got string)
90 }
91
92 type wantExposedConfig struct {
93 discovererType string
94 name string
95 sourceType string
96 status dyncfg.Status
97 }
98
99 func (s *dyncfgSim) run(t *testing.T) {
100 t.Helper()
101
102 require.NotNil(t, s.do, "s.do is nil")
103
104 var buf bytes.Buffer
105 sd := &ServiceDiscovery{
106 Logger: logger.New(),
107 pluginName: testPluginName,
108 dyncfgApi: dyncfg.NewResponder(netdataapi.New(safewriter.New(&buf))),
109 seen: dyncfg.NewSeenCache[sdConfig](),
110 exposed: dyncfg.NewExposedCache[sdConfig](),
111 dyncfgCh: make(chan dyncfg.Function, 1),
112 discoverers: testDiscovererRegistry(),
113 newPipeline: func(cfg pipeline.Config) (sdPipeline, error) {
114 return newTestPipeline(cfg.Name), nil
115 },
116 }
117 sd.sdCb = &sdCallbacks{sd: sd}
118 sd.handler = dyncfg.NewHandler(dyncfg.HandlerOpts[sdConfig]{
119 Logger: sd.Logger,
120 API: sd.dyncfgApi,
121 Seen: sd.seen,
122 Exposed: sd.exposed,
123 Callbacks: sd.sdCb,
124
125 Path: fmt.Sprintf(dyncfgSDPath, testPluginName),
126 EnableFailCode: 422,
127 JobCommands: []dyncfg.Command{
128 dyncfg.CommandSchema,
129 dyncfg.CommandGet,
130 dyncfg.CommandEnable,
131 dyncfg.CommandDisable,
132 dyncfg.CommandUpdate,
133 dyncfg.CommandTest,
134 dyncfg.CommandUserconfig,
135 },
136 })
137
138 done := make(chan struct{})
139 ctx, cancel := context.WithCancel(context.Background())
140
141 // Create output channel (we don't need to capture output for dyncfg tests)
142 out := make(chan<- []*confgroup.Group)
143
144 // Create send function
145 send := func(ctx context.Context, groups []*confgroup.Group) {
146 select {
147 case <-ctx.Done():
148 case out <- groups:
149 }
150 }
151
152 sd.ctx = ctx
153 sd.mgr = NewPipelineManager(sd.Logger, sd.newPipeline, send)
154
155 // Register dyncfg templates (creates CONFIG entries for templates)
156 sd.registerDyncfgTemplates(ctx)
157
158 // Start processing dyncfg commands
159 go func() {
160 defer close(done)
161 for {
162 select {
163 case <-ctx.Done():
164 return
165 case fn := <-sd.dyncfgCh:
166 sd.dyncfgSeqExec(fn)
167 }
168 }
169 }()
170
171 timeout := time.Second * 5
172
173 // Run the test scenario
174 s.do(sd)
175
176 // Give a bit of time for async operations
177 time.Sleep(100 * time.Millisecond)
178
179 cancel()
180
181 select {
182 case <-done:
183 case <-time.After(timeout):
184 t.Errorf("failed to finish work in %s", timeout)
185 }
186
187 // Filter and normalize dyncfg output (same approach as jobmgr sim_test.go)
188 var lines []string
189 for line := range strings.SplitSeq(buf.String(), "\n") {
190 // Skip template CONFIG lines (registered on startup)
191 if strings.HasPrefix(line, "CONFIG") && strings.Contains(line, " template ") {
192 continue
193 }
194 // Remove timestamp from FUNCTION_RESULT_BEGIN
195 if strings.HasPrefix(line, "FUNCTION_RESULT_BEGIN") {
196 parts := strings.Fields(line)
197 line = strings.Join(parts[:len(parts)-1], " ")
198 }
199 lines = append(lines, line)
200 }
201
202 gotDyncfg := strings.TrimSpace(strings.Join(lines, "\n"))
203
204 if s.wantDyncfgFunc != nil {
205 s.wantDyncfgFunc(t, gotDyncfg)
206 } else if s.wantDyncfg != "" {
207 wantDyncfg := strings.TrimSpace(s.wantDyncfg)
208 assert.Equal(t, wantDyncfg, gotDyncfg, "dyncfg commands")
209 }
210
211 // Verify exposed configs
212 if s.wantExposed != nil {
213 wantLen, gotLen := len(s.wantExposed), sd.exposed.Count()
214 require.Equalf(t, wantLen, gotLen, "exposedConfigs: different len (want %d got %d)", wantLen, gotLen)
215
216 for _, want := range s.wantExposed {
217 entry, ok := sd.exposed.LookupByKey(want.discovererType + ":" + want.name)
218 require.Truef(t, ok, "exposedConfigs: config '%s:%s' not found", want.discovererType, want.name)
219 assert.Equal(t, want.sourceType, entry.Cfg.SourceType(), "exposedConfigs: wrong sourceType for '%s:%s'", want.discovererType, want.name)
220 assert.Equal(t, want.status, entry.Status, "exposedConfigs: wrong status for '%s:%s'", want.discovererType, want.name)
221 }
222 }
223
224 // Verify running pipelines
225 if s.wantRunning != nil {
226 gotRunning := sd.mgr.Keys()
227 assert.ElementsMatch(t, s.wantRunning, gotRunning, "running pipelines")
228 }
229 }
230
231 // sendDyncfgCmd sends a dyncfg command and waits for processing
232 func sendDyncfgCmd(sd *ServiceDiscovery, uid string, args []string, payload []byte, source string) {
233 fn := dyncfg.NewFunction(functions.Function{
234 UID: uid,
235 Args: args,
236 Payload: payload,
237 Source: source,
238 ContentType: "application/json",
239 })
240
241 // Call dyncfgConfig directly for commands that are handled there (schema, get, userconfig)
242 // and dyncfgSeqExec for state-changing commands
243 cmd := ""
244 if len(args) >= 2 {
245 cmd = args[1]
246 }
247
248 switch cmd {
249 case "schema", "get", "userconfig", "test":
250 // These are handled directly in dyncfgConfig (read-only/validation commands)
251 sd.dyncfgConfig(fn)
252 default:
253 // State-changing commands go through the channel
254 select {
255 case sd.dyncfgCh <- fn:
256 case <-time.After(time.Second):
257 }
258 }
259
260 // Give time for processing
261 time.Sleep(50 * time.Millisecond)
262 }
263
264 func TestServiceDiscovery_DyncfgSchema(t *testing.T) {
265 tests := map[string]struct {
266 createSim func() *dyncfgSim
267 }{
268 "schema for net_listeners template": {
269 createSim: func() *dyncfgSim {
270 return &dyncfgSim{
271 do: func(sd *ServiceDiscovery) {
272 sendDyncfgCmd(sd, "1-schema",
273 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "schema"},
274 nil, "")
275 },
276 wantDyncfgFunc: func(t *testing.T, got string) {
277 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-schema 200 application/json")
278 assert.Contains(t, got, `"jsonSchema"`)
279 assert.Contains(t, got, "FUNCTION_RESULT_END")
280 },
281 }
282 },
283 },
284 "schema for unknown discoverer type": {
285 createSim: func() *dyncfgSim {
286 return &dyncfgSim{
287 do: func(sd *ServiceDiscovery) {
288 sendDyncfgCmd(sd, "1-schema",
289 []string{sd.dyncfgSDPrefixValue() + "unknown", "schema"},
290 nil, "")
291 },
292 wantDyncfgFunc: func(t *testing.T, got string) {
293 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-schema 404 application/json")
294 assert.Contains(t, got, "Unknown discoverer type")
295 assert.Contains(t, got, "FUNCTION_RESULT_END")
296 },
297 }
298 },
299 },
300 }
301
302 for name, tc := range tests {
303 t.Run(name, func(t *testing.T) {
304 sim := tc.createSim()
305 sim.run(t)
306 })
307 }
308 }
309
310 func TestServiceDiscovery_DyncfgAdd(t *testing.T) {
311 tests := map[string]struct {
312 createSim func() *dyncfgSim
313 }{
314 "add net_listeners job": {
315 createSim: func() *dyncfgSim {
316 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
317 payload, _ := json.Marshal(cfg)
318
319 return &dyncfgSim{
320 do: func(sd *ServiceDiscovery) {
321 sendDyncfgCmd(sd, "1-add",
322 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
323 payload, "type=dyncfg,user=test")
324 },
325 wantExposed: []wantExposedConfig{
326 {
327 discovererType: testDiscovererTypeNetListeners,
328 name: "test-job",
329 sourceType: "dyncfg",
330 status: dyncfg.StatusAccepted,
331 },
332 },
333 wantDyncfg: `
334 FUNCTION_RESULT_BEGIN 1-add 202 application/json
335 {"status":202,"message":""}
336 FUNCTION_RESULT_END
337
338 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
339 `,
340 }
341 },
342 },
343 "add without payload fails": {
344 createSim: func() *dyncfgSim {
345 return &dyncfgSim{
346 do: func(sd *ServiceDiscovery) {
347 sendDyncfgCmd(sd, "1-add",
348 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
349 nil, "type=dyncfg,user=test")
350 },
351 wantExposed: []wantExposedConfig{},
352 wantDyncfg: `
353 FUNCTION_RESULT_BEGIN 1-add 400 application/json
354 {"status":400,"errorMessage":"missing configuration payload"}
355 FUNCTION_RESULT_END
356 `,
357 }
358 },
359 },
360 "add duplicate job replaces existing": {
361 createSim: func() *dyncfgSim {
362 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
363 payload, _ := json.Marshal(cfg)
364
365 return &dyncfgSim{
366 do: func(sd *ServiceDiscovery) {
367 // First add
368 sendDyncfgCmd(sd, "1-add",
369 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
370 payload, "type=dyncfg,user=test")
371
372 // Second add (replaces first - matching jobmgr pattern)
373 sendDyncfgCmd(sd, "2-add",
374 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
375 payload, "type=dyncfg,user=test")
376 },
377 wantExposed: []wantExposedConfig{
378 {
379 discovererType: testDiscovererTypeNetListeners,
380 name: "test-job",
381 sourceType: "dyncfg",
382 status: dyncfg.StatusAccepted,
383 },
384 },
385 wantDyncfg: `
386 FUNCTION_RESULT_BEGIN 1-add 202 application/json
387 {"status":202,"message":""}
388 FUNCTION_RESULT_END
389
390 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
391
392 FUNCTION_RESULT_BEGIN 2-add 202 application/json
393 {"status":202,"message":""}
394 FUNCTION_RESULT_END
395
396 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
397 `,
398 }
399 },
400 },
401 }
402
403 for name, tc := range tests {
404 t.Run(name, func(t *testing.T) {
405 sim := tc.createSim()
406 sim.run(t)
407 })
408 }
409 }
410
411 func TestServiceDiscovery_DyncfgGet(t *testing.T) {
412 tests := map[string]struct {
413 createSim func() *dyncfgSim
414 }{
415 "get existing job": {
416 createSim: func() *dyncfgSim {
417 cfg := newTestNetListenersConfig("serialized-name", 0, 0, defaultTestServices())
418 payload, _ := json.Marshal(cfg)
419
420 return &dyncfgSim{
421 do: func(sd *ServiceDiscovery) {
422 // Add first
423 sendDyncfgCmd(sd, "1-add",
424 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
425 payload, "type=dyncfg,user=test")
426
427 // Get
428 sendDyncfgCmd(sd, "2-get",
429 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "get"},
430 nil, "")
431 },
432 wantDyncfgFunc: func(t *testing.T, got string) {
433 // Check that CONFIG and FUNCTION_RESULT lines are present
434 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create accepted job")
435 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
436 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-get 200 application/json")
437 // JSON key order may vary, so check for presence of expected fields
438 assert.Contains(t, got, `"name":"test-job"`)
439 assert.NotContains(t, got, `"name":"serialized-name"`)
440 assert.Contains(t, got, `"discoverer":{`)
441 assert.Contains(t, got, `"net_listeners":{}`)
442 },
443 }
444 },
445 },
446 "get non-existent job fails": {
447 createSim: func() *dyncfgSim {
448 return &dyncfgSim{
449 do: func(sd *ServiceDiscovery) {
450 sendDyncfgCmd(sd, "1-get",
451 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "non-existent"), "get"},
452 nil, "")
453 },
454 wantDyncfg: `
455 FUNCTION_RESULT_BEGIN 1-get 404 application/json
456 {"status":404,"errorMessage":"Config 'net_listeners:non-existent' not found."}
457 FUNCTION_RESULT_END
458 `,
459 }
460 },
461 },
462 }
463
464 for name, tc := range tests {
465 t.Run(name, func(t *testing.T) {
466 sim := tc.createSim()
467 sim.run(t)
468 })
469 }
470 }
471
472 func TestServiceDiscovery_DyncfgEnableDisable(t *testing.T) {
473 tests := map[string]struct {
474 createSim func() *dyncfgSim
475 }{
476 "enable starts pipeline": {
477 createSim: func() *dyncfgSim {
478 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
479 payload, _ := json.Marshal(cfg)
480
481 return &dyncfgSim{
482 do: func(sd *ServiceDiscovery) {
483 // Add
484 sendDyncfgCmd(sd, "1-add",
485 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
486 payload, "type=dyncfg,user=test")
487
488 // Enable
489 sendDyncfgCmd(sd, "2-enable",
490 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
491 nil, "")
492 },
493 wantExposed: []wantExposedConfig{
494 {
495 discovererType: testDiscovererTypeNetListeners,
496 name: "test-job",
497 sourceType: "dyncfg",
498 status: dyncfg.StatusRunning,
499 },
500 },
501 wantRunning: []string{"dyncfg:net_listeners:test-job"},
502 wantDyncfg: `
503 FUNCTION_RESULT_BEGIN 1-add 202 application/json
504 {"status":202,"message":""}
505 FUNCTION_RESULT_END
506
507 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
508
509 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
510 {"status":200,"message":""}
511 FUNCTION_RESULT_END
512
513 CONFIG test:sd:net_listeners:test-job status running
514 `,
515 }
516 },
517 },
518 "disable stops pipeline": {
519 createSim: func() *dyncfgSim {
520 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
521 payload, _ := json.Marshal(cfg)
522
523 return &dyncfgSim{
524 do: func(sd *ServiceDiscovery) {
525 // Add
526 sendDyncfgCmd(sd, "1-add",
527 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
528 payload, "type=dyncfg,user=test")
529
530 // Enable
531 sendDyncfgCmd(sd, "2-enable",
532 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
533 nil, "")
534
535 // Disable
536 sendDyncfgCmd(sd, "3-disable",
537 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "disable"},
538 nil, "")
539 },
540 wantExposed: []wantExposedConfig{
541 {
542 discovererType: testDiscovererTypeNetListeners,
543 name: "test-job",
544 sourceType: "dyncfg",
545 status: dyncfg.StatusDisabled,
546 },
547 },
548 wantRunning: []string{},
549 wantDyncfg: `
550 FUNCTION_RESULT_BEGIN 1-add 202 application/json
551 {"status":202,"message":""}
552 FUNCTION_RESULT_END
553
554 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
555
556 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
557 {"status":200,"message":""}
558 FUNCTION_RESULT_END
559
560 CONFIG test:sd:net_listeners:test-job status running
561
562 FUNCTION_RESULT_BEGIN 3-disable 200 application/json
563 {"status":200,"message":""}
564 FUNCTION_RESULT_END
565
566 CONFIG test:sd:net_listeners:test-job status disabled
567 `,
568 }
569 },
570 },
571 "enable non-existent job fails": {
572 createSim: func() *dyncfgSim {
573 return &dyncfgSim{
574 do: func(sd *ServiceDiscovery) {
575 sendDyncfgCmd(sd, "1-enable",
576 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "non-existent"), "enable"},
577 nil, "")
578 },
579 wantDyncfg: `
580 FUNCTION_RESULT_BEGIN 1-enable 404 application/json
581 {"status":404,"errorMessage":"job not found."}
582 FUNCTION_RESULT_END
583 `,
584 }
585 },
586 },
587 }
588
589 for name, tc := range tests {
590 t.Run(name, func(t *testing.T) {
591 sim := tc.createSim()
592 sim.run(t)
593 })
594 }
595 }
596
597 func TestServiceDiscovery_DyncfgUpdate(t *testing.T) {
598 tests := map[string]struct {
599 createSim func() *dyncfgSim
600 }{
601 "update disabled job": {
602 createSim: func() *dyncfgSim {
603 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
604 payload, _ := json.Marshal(cfg)
605
606 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
607 updatedPayload, _ := json.Marshal(updatedCfg)
608
609 return &dyncfgSim{
610 do: func(sd *ServiceDiscovery) {
611 // Add
612 sendDyncfgCmd(sd, "1-add",
613 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
614 payload, "type=dyncfg,user=test")
615
616 // Enable then disable to get to Disabled state
617 sendDyncfgCmd(sd, "2-enable",
618 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
619 nil, "")
620
621 sendDyncfgCmd(sd, "3-disable",
622 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "disable"},
623 nil, "")
624
625 // Update (should work in Disabled state)
626 sendDyncfgCmd(sd, "4-update",
627 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
628 updatedPayload, "type=dyncfg,user=test")
629 },
630 wantExposed: []wantExposedConfig{
631 {
632 discovererType: testDiscovererTypeNetListeners,
633 name: "test-job",
634 sourceType: "dyncfg",
635 status: dyncfg.StatusDisabled,
636 },
637 },
638 wantDyncfg: `
639 FUNCTION_RESULT_BEGIN 1-add 202 application/json
640 {"status":202,"message":""}
641 FUNCTION_RESULT_END
642
643 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
644
645 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
646 {"status":200,"message":""}
647 FUNCTION_RESULT_END
648
649 CONFIG test:sd:net_listeners:test-job status running
650
651 FUNCTION_RESULT_BEGIN 3-disable 200 application/json
652 {"status":200,"message":""}
653 FUNCTION_RESULT_END
654
655 CONFIG test:sd:net_listeners:test-job status disabled
656
657 FUNCTION_RESULT_BEGIN 4-update 200 application/json
658 {"status":200,"message":""}
659 FUNCTION_RESULT_END
660
661 CONFIG test:sd:net_listeners:test-job status disabled
662 `,
663 }
664 },
665 },
666 "update in accepted state fails": {
667 createSim: func() *dyncfgSim {
668 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
669 payload, _ := json.Marshal(cfg)
670
671 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
672 updatedPayload, _ := json.Marshal(updatedCfg)
673
674 return &dyncfgSim{
675 do: func(sd *ServiceDiscovery) {
676 // Add (creates in Accepted state)
677 sendDyncfgCmd(sd, "1-add",
678 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
679 payload, "type=dyncfg,user=test")
680
681 // Update in Accepted state should fail with 403
682 sendDyncfgCmd(sd, "2-update",
683 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
684 updatedPayload, "type=dyncfg,user=test")
685 },
686 wantExposed: []wantExposedConfig{
687 {
688 discovererType: testDiscovererTypeNetListeners,
689 name: "test-job",
690 sourceType: "dyncfg",
691 status: dyncfg.StatusAccepted,
692 },
693 },
694 wantDyncfg: `
695 FUNCTION_RESULT_BEGIN 1-add 202 application/json
696 {"status":202,"message":""}
697 FUNCTION_RESULT_END
698
699 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
700
701 FUNCTION_RESULT_BEGIN 2-update 403 application/json
702 {"status":403,"errorMessage":"updating is not allowed in 'accepted' state."}
703 FUNCTION_RESULT_END
704
705 CONFIG test:sd:net_listeners:test-job status accepted
706 `,
707 }
708 },
709 },
710 "update non-existent job fails": {
711 createSim: func() *dyncfgSim {
712 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
713 payload, _ := json.Marshal(cfg)
714
715 return &dyncfgSim{
716 do: func(sd *ServiceDiscovery) {
717 sendDyncfgCmd(sd, "1-update",
718 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "non-existent"), "update"},
719 payload, "type=dyncfg,user=test")
720 },
721 wantDyncfg: `
722 FUNCTION_RESULT_BEGIN 1-update 404 application/json
723 {"status":404,"errorMessage":"job not found."}
724 FUNCTION_RESULT_END
725 `,
726 }
727 },
728 },
729 }
730
731 for name, tc := range tests {
732 t.Run(name, func(t *testing.T) {
733 sim := tc.createSim()
734 sim.run(t)
735 })
736 }
737 }
738
739 func TestServiceDiscovery_DyncfgRemove(t *testing.T) {
740 tests := map[string]struct {
741 createSim func() *dyncfgSim
742 }{
743 "remove dyncfg job": {
744 createSim: func() *dyncfgSim {
745 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
746 payload, _ := json.Marshal(cfg)
747
748 return &dyncfgSim{
749 do: func(sd *ServiceDiscovery) {
750 // Add
751 sendDyncfgCmd(sd, "1-add",
752 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
753 payload, "type=dyncfg,user=test")
754
755 // Remove
756 sendDyncfgCmd(sd, "2-remove",
757 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "remove"},
758 nil, "")
759 },
760 wantExposed: []wantExposedConfig{},
761 wantRunning: []string{},
762 wantDyncfg: `
763 FUNCTION_RESULT_BEGIN 1-add 202 application/json
764 {"status":202,"message":""}
765 FUNCTION_RESULT_END
766
767 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
768
769 FUNCTION_RESULT_BEGIN 2-remove 200 application/json
770 {"status":200,"message":""}
771 FUNCTION_RESULT_END
772
773 CONFIG test:sd:net_listeners:test-job delete
774 `,
775 }
776 },
777 },
778 "remove running job stops it first": {
779 createSim: func() *dyncfgSim {
780 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
781 payload, _ := json.Marshal(cfg)
782
783 return &dyncfgSim{
784 do: func(sd *ServiceDiscovery) {
785 // Add
786 sendDyncfgCmd(sd, "1-add",
787 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
788 payload, "type=dyncfg,user=test")
789
790 // Enable
791 sendDyncfgCmd(sd, "2-enable",
792 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
793 nil, "")
794
795 // Remove (should stop first)
796 sendDyncfgCmd(sd, "3-remove",
797 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "remove"},
798 nil, "")
799 },
800 wantExposed: []wantExposedConfig{},
801 wantRunning: []string{},
802 wantDyncfg: `
803 FUNCTION_RESULT_BEGIN 1-add 202 application/json
804 {"status":202,"message":""}
805 FUNCTION_RESULT_END
806
807 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
808
809 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
810 {"status":200,"message":""}
811 FUNCTION_RESULT_END
812
813 CONFIG test:sd:net_listeners:test-job status running
814
815 FUNCTION_RESULT_BEGIN 3-remove 200 application/json
816 {"status":200,"message":""}
817 FUNCTION_RESULT_END
818
819 CONFIG test:sd:net_listeners:test-job delete
820 `,
821 }
822 },
823 },
824 "remove non-existent job fails": {
825 createSim: func() *dyncfgSim {
826 return &dyncfgSim{
827 do: func(sd *ServiceDiscovery) {
828 sendDyncfgCmd(sd, "1-remove",
829 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "non-existent"), "remove"},
830 nil, "")
831 },
832 wantDyncfg: `
833 FUNCTION_RESULT_BEGIN 1-remove 404 application/json
834 {"status":404,"errorMessage":"job not found."}
835 FUNCTION_RESULT_END
836 `,
837 }
838 },
839 },
840 }
841
842 for name, tc := range tests {
843 t.Run(name, func(t *testing.T) {
844 sim := tc.createSim()
845 sim.run(t)
846 })
847 }
848 }
849
850 func TestServiceDiscovery_DyncfgUserconfig(t *testing.T) {
851 tests := map[string]struct {
852 createSim func() *dyncfgSim
853 }{
854 "userconfig for template": {
855 createSim: func() *dyncfgSim {
856 cfg := newTestNetListenersConfig("serialized-name", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
857 payload, _ := json.Marshal(cfg)
858
859 return &dyncfgSim{
860 do: func(sd *ServiceDiscovery) {
861 sendDyncfgCmd(sd, "1-userconfig",
862 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "userconfig"},
863 payload, "")
864 },
865 wantDyncfgFunc: func(t *testing.T, got string) {
866 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-userconfig 200 application/yaml")
867 assert.Contains(t, got, "name: test")
868 assert.NotContains(t, got, "name: serialized-name")
869 assert.Contains(t, got, "discoverer:")
870 assert.Contains(t, got, "net_listeners:")
871 assert.Contains(t, got, "interval: 5")
872 assert.Contains(t, got, "FUNCTION_RESULT_END")
873 },
874 }
875 },
876 },
877 "userconfig for existing job": {
878 createSim: func() *dyncfgSim {
879 cfg := newTestNetListenersConfig("serialized-name", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
880 payload, _ := json.Marshal(cfg)
881
882 return &dyncfgSim{
883 do: func(sd *ServiceDiscovery) {
884 // Add first
885 sendDyncfgCmd(sd, "1-add",
886 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
887 payload, "type=dyncfg,user=test")
888
889 // Userconfig - must provide payload (matching jobmgr pattern)
890 sendDyncfgCmd(sd, "2-userconfig",
891 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "userconfig"},
892 payload, "")
893 },
894 wantDyncfgFunc: func(t *testing.T, got string) {
895 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
896 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-userconfig 200 application/yaml")
897 assert.Contains(t, got, "name: test-job")
898 assert.NotContains(t, got, "name: serialized-name")
899 assert.Contains(t, got, "discoverer:")
900 assert.Contains(t, got, "net_listeners:")
901 assert.Contains(t, got, "interval: 5")
902 },
903 }
904 },
905 },
906 "userconfig for docker template": {
907 createSim: func() *dyncfgSim {
908 cfg := newTestDockerConfig("docker-test", "unix:///var/run/docker.sock", confopt.Duration(5*time.Second), defaultTestServices())
909 payload, _ := json.Marshal(cfg)
910
911 return &dyncfgSim{
912 do: func(sd *ServiceDiscovery) {
913 sendDyncfgCmd(sd, "1-userconfig",
914 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "userconfig"},
915 payload, "")
916 },
917 wantDyncfgFunc: func(t *testing.T, got string) {
918 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-userconfig 200 application/yaml")
919 assert.Contains(t, got, "discoverer:")
920 assert.Contains(t, got, "docker:")
921 assert.Contains(t, got, "address: unix:///var/run/docker.sock")
922 },
923 }
924 },
925 },
926 "userconfig for k8s template": {
927 createSim: func() *dyncfgSim {
928 cfg := newTestK8sConfig("k8s-test", []testK8sConfig{{
929 Role: "pod",
930 Namespaces: []string{"default"},
931 }}, defaultTestServices())
932 payload, _ := json.Marshal(cfg)
933
934 return &dyncfgSim{
935 do: func(sd *ServiceDiscovery) {
936 sendDyncfgCmd(sd, "1-userconfig",
937 []string{sd.dyncfgTemplateID(testDiscovererTypeK8s), "userconfig"},
938 payload, "")
939 },
940 wantDyncfgFunc: func(t *testing.T, got string) {
941 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-userconfig 200 application/yaml")
942 assert.Contains(t, got, "discoverer:")
943 assert.Contains(t, got, "k8s:")
944 assert.Contains(t, got, "role: pod")
945 assert.Contains(t, got, "- default")
946 },
947 }
948 },
949 },
950 "userconfig for snmp template": {
951 createSim: func() *dyncfgSim {
952 cfg := newTestSNMPConfig("snmp-test", testSNMPConfig{
953 RescanInterval: confopt.LongDuration(1 * time.Hour),
954 Credentials: []testSNMPCredentialConfig{{Name: "public-v2", Version: "2c", Community: "public"}},
955 Networks: []testSNMPNetworkConfig{{Subnet: "192.168.1.0/24", Credential: "public-v2"}},
956 }, defaultTestServices())
957 payload, _ := json.Marshal(cfg)
958
959 return &dyncfgSim{
960 do: func(sd *ServiceDiscovery) {
961 sendDyncfgCmd(sd, "1-userconfig",
962 []string{sd.dyncfgTemplateID(testDiscovererTypeSNMP), "userconfig"},
963 payload, "")
964 },
965 wantDyncfgFunc: func(t *testing.T, got string) {
966 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-userconfig 200 application/yaml")
967 assert.Contains(t, got, "discoverer:")
968 assert.Contains(t, got, "snmp:")
969 assert.Contains(t, got, "rescan_interval: 1h")
970 assert.Contains(t, got, "subnet: 192.168.1.0/24")
971 },
972 }
973 },
974 },
975 "userconfig fails on mismatched discoverer type": {
976 createSim: func() *dyncfgSim {
977 cfg := newTestDockerConfig("docker-test", "unix:///var/run/docker.sock", confopt.Duration(5*time.Second), defaultTestServices())
978 payload, _ := json.Marshal(cfg)
979
980 return &dyncfgSim{
981 do: func(sd *ServiceDiscovery) {
982 sendDyncfgCmd(sd, "1-userconfig",
983 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "userconfig"},
984 payload, "")
985 },
986 wantDyncfgFunc: func(t *testing.T, got string) {
987 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-userconfig 400 application/json")
988 assert.Contains(t, got, "Failed to parse config")
989 assert.Contains(t, got, "expected")
990 },
991 }
992 },
993 },
994 }
995
996 for name, tc := range tests {
997 t.Run(name, func(t *testing.T) {
998 sim := tc.createSim()
999 sim.run(t)
1000 })
1001 }
1002 }
1003
1004 func TestServiceDiscovery_DyncfgFileConfig(t *testing.T) {
1005 tests := map[string]struct {
1006 createSim func() *dyncfgSim
1007 }{
1008 "file config cannot be removed via dyncfg": {
1009 createSim: func() *dyncfgSim {
1010 return &dyncfgSim{
1011 do: func(sd *ServiceDiscovery) {
1012 // Manually add a file-based config to exposed cache
1013 cfg := sdConfig{
1014 "name": "file-config",
1015 ikeyDiscovererType: testDiscovererTypeNetListeners,
1016 ikeyPipelineKey: "/etc/netdata/sd/test.conf",
1017 ikeySource: "/etc/netdata/sd/test.conf",
1018 ikeySourceType: "file",
1019 }
1020 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: cfg, Status: dyncfg.StatusRunning})
1021
1022 // Try to remove
1023 sendDyncfgCmd(sd, "1-remove",
1024 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "file-config"), "remove"},
1025 nil, "")
1026 },
1027 wantExposed: []wantExposedConfig{
1028 {
1029 discovererType: testDiscovererTypeNetListeners,
1030 name: "file-config",
1031 sourceType: "file",
1032 status: dyncfg.StatusRunning,
1033 },
1034 },
1035 wantDyncfg: `
1036 FUNCTION_RESULT_BEGIN 1-remove 405 application/json
1037 {"status":405,"errorMessage":"removing jobs of type 'file' is not supported, only 'dyncfg' jobs can be removed."}
1038 FUNCTION_RESULT_END
1039 `,
1040 }
1041 },
1042 },
1043 }
1044
1045 for name, tc := range tests {
1046 t.Run(name, func(t *testing.T) {
1047 sim := tc.createSim()
1048 sim.run(t)
1049 })
1050 }
1051 }
1052
1053 func TestServiceDiscovery_DyncfgDockerConfig(t *testing.T) {
1054 tests := map[string]struct {
1055 createSim func() *dyncfgSim
1056 }{
1057 "add docker job": {
1058 createSim: func() *dyncfgSim {
1059 cfg := newTestDockerConfig("docker-test", "unix:///var/run/docker.sock", confopt.Duration(5*time.Second), []pipeline.ServiceRuleConfig{
1060 {ID: "nginx", Match: `{{ glob .Image "*nginx*" }}`},
1061 })
1062 payload, _ := json.Marshal(cfg)
1063
1064 return &dyncfgSim{
1065 do: func(sd *ServiceDiscovery) {
1066 sendDyncfgCmd(sd, "1-add",
1067 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "add", "docker-test"},
1068 payload, "type=dyncfg,user=test")
1069 },
1070 wantExposed: []wantExposedConfig{
1071 {
1072 discovererType: testDiscovererTypeDocker,
1073 name: "docker-test",
1074 sourceType: "dyncfg",
1075 status: dyncfg.StatusAccepted,
1076 },
1077 },
1078 wantDyncfg: `
1079 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1080 {"status":202,"message":""}
1081 FUNCTION_RESULT_END
1082
1083 CONFIG test:sd:docker:docker-test create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1084 `,
1085 }
1086 },
1087 },
1088 "add and enable docker job": {
1089 createSim: func() *dyncfgSim {
1090 cfg := newTestDockerConfig("docker-test", "tcp://localhost:2375", 0, defaultTestServices())
1091 payload, _ := json.Marshal(cfg)
1092
1093 return &dyncfgSim{
1094 do: func(sd *ServiceDiscovery) {
1095 // Add
1096 sendDyncfgCmd(sd, "1-add",
1097 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "add", "docker-test"},
1098 payload, "type=dyncfg,user=test")
1099
1100 // Enable
1101 sendDyncfgCmd(sd, "2-enable",
1102 []string{sd.dyncfgJobID(testDiscovererTypeDocker, "docker-test"), "enable"},
1103 nil, "")
1104 },
1105 wantExposed: []wantExposedConfig{
1106 {
1107 discovererType: testDiscovererTypeDocker,
1108 name: "docker-test",
1109 sourceType: "dyncfg",
1110 status: dyncfg.StatusRunning,
1111 },
1112 },
1113 wantRunning: []string{"dyncfg:docker:docker-test"},
1114 }
1115 },
1116 },
1117 "get docker job config": {
1118 createSim: func() *dyncfgSim {
1119 cfg := newTestDockerConfig("docker-test", "unix:///var/run/docker.sock", 0, defaultTestServices())
1120 payload, _ := json.Marshal(cfg)
1121
1122 return &dyncfgSim{
1123 do: func(sd *ServiceDiscovery) {
1124 // Add
1125 sendDyncfgCmd(sd, "1-add",
1126 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "add", "docker-test"},
1127 payload, "type=dyncfg,user=test")
1128
1129 // Get
1130 sendDyncfgCmd(sd, "2-get",
1131 []string{sd.dyncfgJobID(testDiscovererTypeDocker, "docker-test"), "get"},
1132 nil, "")
1133 },
1134 wantDyncfgFunc: func(t *testing.T, got string) {
1135 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-get 200 application/json")
1136 assert.Contains(t, got, `"discoverer":{`)
1137 assert.Contains(t, got, `"docker":{`)
1138 assert.Contains(t, got, `"address":"unix:///var/run/docker.sock"`)
1139 },
1140 }
1141 },
1142 },
1143 }
1144
1145 for name, tc := range tests {
1146 t.Run(name, func(t *testing.T) {
1147 sim := tc.createSim()
1148 sim.run(t)
1149 })
1150 }
1151 }
1152
1153 func TestServiceDiscovery_DyncfgK8sConfig(t *testing.T) {
1154 tests := map[string]struct {
1155 createSim func() *dyncfgSim
1156 }{
1157 "add k8s job": {
1158 createSim: func() *dyncfgSim {
1159 k8sCfg := testK8sConfig{
1160 Role: "pod",
1161 Namespaces: []string{"default", "kube-system"},
1162 }
1163 k8sCfg.Selector.Label = "app=nginx"
1164 k8sCfg.Pod.LocalMode = true
1165 cfg := newTestK8sConfig("k8s-test", []testK8sConfig{k8sCfg}, []pipeline.ServiceRuleConfig{
1166 {ID: "nginx-pods", Match: `{{ eq .Namespace "default" }}`},
1167 })
1168 payload, _ := json.Marshal(cfg)
1169
1170 return &dyncfgSim{
1171 do: func(sd *ServiceDiscovery) {
1172 sendDyncfgCmd(sd, "1-add",
1173 []string{sd.dyncfgTemplateID(testDiscovererTypeK8s), "add", "k8s-test"},
1174 payload, "type=dyncfg,user=test")
1175 },
1176 wantExposed: []wantExposedConfig{
1177 {
1178 discovererType: testDiscovererTypeK8s,
1179 name: "k8s-test",
1180 sourceType: "dyncfg",
1181 status: dyncfg.StatusAccepted,
1182 },
1183 },
1184 wantDyncfg: `
1185 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1186 {"status":202,"message":""}
1187 FUNCTION_RESULT_END
1188
1189 CONFIG test:sd:k8s:k8s-test create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1190 `,
1191 }
1192 },
1193 },
1194 "add k8s service role job": {
1195 createSim: func() *dyncfgSim {
1196 cfg := newTestK8sConfig("k8s-svc-test", []testK8sConfig{{Role: "service"}}, defaultTestServices())
1197 payload, _ := json.Marshal(cfg)
1198
1199 return &dyncfgSim{
1200 do: func(sd *ServiceDiscovery) {
1201 // Add
1202 sendDyncfgCmd(sd, "1-add",
1203 []string{sd.dyncfgTemplateID(testDiscovererTypeK8s), "add", "k8s-svc-test"},
1204 payload, "type=dyncfg,user=test")
1205
1206 // Enable
1207 sendDyncfgCmd(sd, "2-enable",
1208 []string{sd.dyncfgJobID(testDiscovererTypeK8s, "k8s-svc-test"), "enable"},
1209 nil, "")
1210 },
1211 wantExposed: []wantExposedConfig{
1212 {
1213 discovererType: testDiscovererTypeK8s,
1214 name: "k8s-svc-test",
1215 sourceType: "dyncfg",
1216 status: dyncfg.StatusRunning,
1217 },
1218 },
1219 wantRunning: []string{"dyncfg:k8s:k8s-svc-test"},
1220 }
1221 },
1222 },
1223 "get k8s job config": {
1224 createSim: func() *dyncfgSim {
1225 cfg := newTestK8sConfig("k8s-test", []testK8sConfig{{Role: "pod", Namespaces: []string{"default"}}}, defaultTestServices())
1226 payload, _ := json.Marshal(cfg)
1227
1228 return &dyncfgSim{
1229 do: func(sd *ServiceDiscovery) {
1230 // Add
1231 sendDyncfgCmd(sd, "1-add",
1232 []string{sd.dyncfgTemplateID(testDiscovererTypeK8s), "add", "k8s-test"},
1233 payload, "type=dyncfg,user=test")
1234
1235 // Get
1236 sendDyncfgCmd(sd, "2-get",
1237 []string{sd.dyncfgJobID(testDiscovererTypeK8s, "k8s-test"), "get"},
1238 nil, "")
1239 },
1240 wantDyncfgFunc: func(t *testing.T, got string) {
1241 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-get 200 application/json")
1242 assert.Contains(t, got, `"discoverer":{`)
1243 assert.Contains(t, got, `"k8s":[`)
1244 assert.Contains(t, got, `"role":"pod"`)
1245 assert.Contains(t, got, `"namespaces":["default"]`)
1246 },
1247 }
1248 },
1249 },
1250 }
1251
1252 for name, tc := range tests {
1253 t.Run(name, func(t *testing.T) {
1254 sim := tc.createSim()
1255 sim.run(t)
1256 })
1257 }
1258 }
1259
1260 func TestServiceDiscovery_DyncfgSNMPConfig(t *testing.T) {
1261 tests := map[string]struct {
1262 createSim func() *dyncfgSim
1263 }{
1264 "add snmp job": {
1265 createSim: func() *dyncfgSim {
1266 cfg := newTestSNMPConfig("snmp-test", testSNMPConfig{
1267 RescanInterval: confopt.LongDuration(30 * time.Minute),
1268 Timeout: confopt.Duration(1 * time.Second),
1269 DeviceCacheTTL: confopt.LongDuration(12 * time.Hour),
1270 Credentials: []testSNMPCredentialConfig{{Name: "public-v2", Version: "2c", Community: "public"}},
1271 Networks: []testSNMPNetworkConfig{{Subnet: "192.168.1.0/24", Credential: "public-v2"}},
1272 }, defaultTestServices())
1273 payload, _ := json.Marshal(cfg)
1274
1275 return &dyncfgSim{
1276 do: func(sd *ServiceDiscovery) {
1277 sendDyncfgCmd(sd, "1-add",
1278 []string{sd.dyncfgTemplateID(testDiscovererTypeSNMP), "add", "snmp-test"},
1279 payload, "type=dyncfg,user=test")
1280 },
1281 wantExposed: []wantExposedConfig{
1282 {
1283 discovererType: testDiscovererTypeSNMP,
1284 name: "snmp-test",
1285 sourceType: "dyncfg",
1286 status: dyncfg.StatusAccepted,
1287 },
1288 },
1289 wantDyncfg: `
1290 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1291 {"status":202,"message":""}
1292 FUNCTION_RESULT_END
1293
1294 CONFIG test:sd:snmp:snmp-test create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1295 `,
1296 }
1297 },
1298 },
1299 "add snmp v3 job": {
1300 createSim: func() *dyncfgSim {
1301 cfg := newTestSNMPConfig("snmp-v3-test", testSNMPConfig{
1302 Credentials: []testSNMPCredentialConfig{{
1303 Name: "snmpv3-auth",
1304 Version: "3",
1305 UserName: "admin",
1306 SecurityLevel: "authPriv",
1307 AuthProtocol: "sha256",
1308 AuthPassphrase: "authpass",
1309 PrivacyProtocol: "aes",
1310 PrivacyPassphrase: "privpass",
1311 }},
1312 Networks: []testSNMPNetworkConfig{{Subnet: "10.0.0.0/24", Credential: "snmpv3-auth"}},
1313 }, defaultTestServices())
1314 payload, _ := json.Marshal(cfg)
1315
1316 return &dyncfgSim{
1317 do: func(sd *ServiceDiscovery) {
1318 // Add
1319 sendDyncfgCmd(sd, "1-add",
1320 []string{sd.dyncfgTemplateID(testDiscovererTypeSNMP), "add", "snmp-v3-test"},
1321 payload, "type=dyncfg,user=test")
1322
1323 // Enable
1324 sendDyncfgCmd(sd, "2-enable",
1325 []string{sd.dyncfgJobID(testDiscovererTypeSNMP, "snmp-v3-test"), "enable"},
1326 nil, "")
1327 },
1328 wantExposed: []wantExposedConfig{
1329 {
1330 discovererType: testDiscovererTypeSNMP,
1331 name: "snmp-v3-test",
1332 sourceType: "dyncfg",
1333 status: dyncfg.StatusRunning,
1334 },
1335 },
1336 wantRunning: []string{"dyncfg:snmp:snmp-v3-test"},
1337 }
1338 },
1339 },
1340 "get snmp job config": {
1341 createSim: func() *dyncfgSim {
1342 cfg := newTestSNMPConfig("snmp-test", testSNMPConfig{
1343 RescanInterval: confopt.LongDuration(1 * time.Hour),
1344 Credentials: []testSNMPCredentialConfig{{Name: "v2-cred", Version: "2c", Community: "public"}},
1345 Networks: []testSNMPNetworkConfig{{Subnet: "192.168.0.0/16", Credential: "v2-cred"}},
1346 }, defaultTestServices())
1347 payload, _ := json.Marshal(cfg)
1348
1349 return &dyncfgSim{
1350 do: func(sd *ServiceDiscovery) {
1351 // Add
1352 sendDyncfgCmd(sd, "1-add",
1353 []string{sd.dyncfgTemplateID(testDiscovererTypeSNMP), "add", "snmp-test"},
1354 payload, "type=dyncfg,user=test")
1355
1356 // Get
1357 sendDyncfgCmd(sd, "2-get",
1358 []string{sd.dyncfgJobID(testDiscovererTypeSNMP, "snmp-test"), "get"},
1359 nil, "")
1360 },
1361 wantDyncfgFunc: func(t *testing.T, got string) {
1362 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-get 200 application/json")
1363 assert.Contains(t, got, `"discoverer":{`)
1364 assert.Contains(t, got, `"snmp":{`)
1365 assert.Contains(t, got, `"rescan_interval":"1h"`)
1366 assert.Contains(t, got, `"subnet":"192.168.0.0/16"`)
1367 },
1368 }
1369 },
1370 },
1371 }
1372
1373 for name, tc := range tests {
1374 t.Run(name, func(t *testing.T) {
1375 sim := tc.createSim()
1376 sim.run(t)
1377 })
1378 }
1379 }
1380
1381 func TestServiceDiscovery_DyncfgUpdateWhileRunning(t *testing.T) {
1382 tests := map[string]struct {
1383 createSim func() *dyncfgSim
1384 }{
1385 "update running pipeline restarts it": {
1386 createSim: func() *dyncfgSim {
1387 cfg := newTestNetListenersConfig("test-job", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
1388 payload, _ := json.Marshal(cfg)
1389
1390 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
1391 updatedPayload, _ := json.Marshal(updatedCfg)
1392
1393 return &dyncfgSim{
1394 do: func(sd *ServiceDiscovery) {
1395 // Add
1396 sendDyncfgCmd(sd, "1-add",
1397 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1398 payload, "type=dyncfg,user=test")
1399
1400 // Enable (starts pipeline)
1401 sendDyncfgCmd(sd, "2-enable",
1402 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
1403 nil, "")
1404
1405 // Update while running (should restart pipeline)
1406 sendDyncfgCmd(sd, "3-update",
1407 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
1408 updatedPayload, "type=dyncfg,user=test")
1409 },
1410 wantExposed: []wantExposedConfig{
1411 {
1412 discovererType: testDiscovererTypeNetListeners,
1413 name: "test-job",
1414 sourceType: "dyncfg",
1415 status: dyncfg.StatusRunning,
1416 },
1417 },
1418 wantRunning: []string{"dyncfg:net_listeners:test-job"},
1419 wantDyncfg: `
1420 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1421 {"status":202,"message":""}
1422 FUNCTION_RESULT_END
1423
1424 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1425
1426 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
1427 {"status":200,"message":""}
1428 FUNCTION_RESULT_END
1429
1430 CONFIG test:sd:net_listeners:test-job status running
1431
1432 FUNCTION_RESULT_BEGIN 3-update 200 application/json
1433 {"status":200,"message":""}
1434 FUNCTION_RESULT_END
1435
1436 CONFIG test:sd:net_listeners:test-job status running
1437 `,
1438 }
1439 },
1440 },
1441 "update running docker pipeline": {
1442 createSim: func() *dyncfgSim {
1443 cfg := newTestDockerConfig("docker-job", "unix:///var/run/docker.sock", 0, defaultTestServices())
1444 payload, _ := json.Marshal(cfg)
1445
1446 updatedCfg := newTestDockerConfig("docker-job", "tcp://localhost:2375", 0, defaultTestServices())
1447 updatedPayload, _ := json.Marshal(updatedCfg)
1448
1449 return &dyncfgSim{
1450 do: func(sd *ServiceDiscovery) {
1451 // Add
1452 sendDyncfgCmd(sd, "1-add",
1453 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "add", "docker-job"},
1454 payload, "type=dyncfg,user=test")
1455
1456 // Enable
1457 sendDyncfgCmd(sd, "2-enable",
1458 []string{sd.dyncfgJobID(testDiscovererTypeDocker, "docker-job"), "enable"},
1459 nil, "")
1460
1461 // Update while running
1462 sendDyncfgCmd(sd, "3-update",
1463 []string{sd.dyncfgJobID(testDiscovererTypeDocker, "docker-job"), "update"},
1464 updatedPayload, "type=dyncfg,user=test")
1465
1466 // Verify config was updated
1467 sendDyncfgCmd(sd, "4-get",
1468 []string{sd.dyncfgJobID(testDiscovererTypeDocker, "docker-job"), "get"},
1469 nil, "")
1470 },
1471 wantExposed: []wantExposedConfig{
1472 {
1473 discovererType: testDiscovererTypeDocker,
1474 name: "docker-job",
1475 sourceType: "dyncfg",
1476 status: dyncfg.StatusRunning,
1477 },
1478 },
1479 wantRunning: []string{"dyncfg:docker:docker-job"},
1480 wantDyncfgFunc: func(t *testing.T, got string) {
1481 // Verify update response
1482 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 3-update 200 application/json")
1483 // Verify config was updated to new address
1484 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 4-get 200 application/json")
1485 assert.Contains(t, got, `"address":"tcp://localhost:2375"`)
1486 },
1487 }
1488 },
1489 },
1490 }
1491
1492 for name, tc := range tests {
1493 t.Run(name, func(t *testing.T) {
1494 sim := tc.createSim()
1495 sim.run(t)
1496 })
1497 }
1498 }
1499
1500 func TestServiceDiscovery_DyncfgTest(t *testing.T) {
1501 tests := map[string]struct {
1502 createSim func() *dyncfgSim
1503 }{
1504 "test valid config succeeds": {
1505 createSim: func() *dyncfgSim {
1506 cfg := newTestNetListenersConfig("test-job", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
1507 payload, _ := json.Marshal(cfg)
1508
1509 return &dyncfgSim{
1510 do: func(sd *ServiceDiscovery) {
1511 sendDyncfgCmd(sd, "1-test",
1512 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "test"},
1513 payload, "")
1514 },
1515 wantExposed: []wantExposedConfig{},
1516 wantRunning: []string{},
1517 wantDyncfg: `
1518 FUNCTION_RESULT_BEGIN 1-test 200 application/json
1519 {"status":200,"message":""}
1520 FUNCTION_RESULT_END
1521 `,
1522 }
1523 },
1524 },
1525 "test without payload fails": {
1526 createSim: func() *dyncfgSim {
1527 return &dyncfgSim{
1528 do: func(sd *ServiceDiscovery) {
1529 sendDyncfgCmd(sd, "1-test",
1530 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "test"},
1531 nil, "")
1532 },
1533 wantExposed: []wantExposedConfig{},
1534 wantRunning: []string{},
1535 wantDyncfg: `
1536 FUNCTION_RESULT_BEGIN 1-test 400 application/json
1537 {"status":400,"errorMessage":"missing configuration payload"}
1538 FUNCTION_RESULT_END
1539 `,
1540 }
1541 },
1542 },
1543 "test valid docker config": {
1544 createSim: func() *dyncfgSim {
1545 cfg := newTestDockerConfig("docker-test", "unix:///var/run/docker.sock", 0, defaultTestServices())
1546 payload, _ := json.Marshal(cfg)
1547
1548 return &dyncfgSim{
1549 do: func(sd *ServiceDiscovery) {
1550 sendDyncfgCmd(sd, "1-test",
1551 []string{sd.dyncfgTemplateID(testDiscovererTypeDocker), "test"},
1552 payload, "")
1553 },
1554 wantExposed: []wantExposedConfig{},
1555 wantRunning: []string{},
1556 wantDyncfg: `
1557 FUNCTION_RESULT_BEGIN 1-test 200 application/json
1558 {"status":200,"message":""}
1559 FUNCTION_RESULT_END
1560 `,
1561 }
1562 },
1563 },
1564 "test valid k8s config": {
1565 createSim: func() *dyncfgSim {
1566 cfg := newTestK8sConfig("k8s-test", []testK8sConfig{{Role: "pod"}}, defaultTestServices())
1567 payload, _ := json.Marshal(cfg)
1568
1569 return &dyncfgSim{
1570 do: func(sd *ServiceDiscovery) {
1571 sendDyncfgCmd(sd, "1-test",
1572 []string{sd.dyncfgTemplateID(testDiscovererTypeK8s), "test"},
1573 payload, "")
1574 },
1575 wantExposed: []wantExposedConfig{},
1576 wantRunning: []string{},
1577 wantDyncfg: `
1578 FUNCTION_RESULT_BEGIN 1-test 200 application/json
1579 {"status":200,"message":""}
1580 FUNCTION_RESULT_END
1581 `,
1582 }
1583 },
1584 },
1585 "test valid snmp config": {
1586 createSim: func() *dyncfgSim {
1587 cfg := newTestSNMPConfig("snmp-test", testSNMPConfig{
1588 Credentials: []testSNMPCredentialConfig{{Name: "v2-cred", Version: "2c"}},
1589 Networks: []testSNMPNetworkConfig{{Subnet: "192.168.1.0/24", Credential: "v2-cred"}},
1590 }, defaultTestServices())
1591 payload, _ := json.Marshal(cfg)
1592
1593 return &dyncfgSim{
1594 do: func(sd *ServiceDiscovery) {
1595 sendDyncfgCmd(sd, "1-test",
1596 []string{sd.dyncfgTemplateID(testDiscovererTypeSNMP), "test"},
1597 payload, "")
1598 },
1599 wantExposed: []wantExposedConfig{},
1600 wantRunning: []string{},
1601 wantDyncfg: `
1602 FUNCTION_RESULT_BEGIN 1-test 200 application/json
1603 {"status":200,"message":""}
1604 FUNCTION_RESULT_END
1605 `,
1606 }
1607 },
1608 },
1609 "test existing job config succeeds": {
1610 createSim: func() *dyncfgSim {
1611 cfg := newTestNetListenersConfig("test-job", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
1612 payload, _ := json.Marshal(cfg)
1613
1614 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
1615 updatedPayload, _ := json.Marshal(updatedCfg)
1616
1617 return &dyncfgSim{
1618 do: func(sd *ServiceDiscovery) {
1619 // First add the job
1620 sendDyncfgCmd(sd, "1-add",
1621 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1622 payload, "type=dyncfg,user=test")
1623
1624 // Test with new config (validates without applying)
1625 sendDyncfgCmd(sd, "2-test",
1626 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "test"},
1627 updatedPayload, "")
1628 },
1629 wantExposed: []wantExposedConfig{
1630 {
1631 discovererType: testDiscovererTypeNetListeners,
1632 name: "test-job",
1633 sourceType: "dyncfg",
1634 status: dyncfg.StatusAccepted, // Still accepted, not changed by test
1635 },
1636 },
1637 wantDyncfgFunc: func(t *testing.T, got string) {
1638 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
1639 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-test 200 application/json")
1640 },
1641 }
1642 },
1643 },
1644 "test job with invalid config fails": {
1645 createSim: func() *dyncfgSim {
1646 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
1647 payload, _ := json.Marshal(cfg)
1648
1649 return &dyncfgSim{
1650 do: func(sd *ServiceDiscovery) {
1651 // First add the job
1652 sendDyncfgCmd(sd, "1-add",
1653 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1654 payload, "type=dyncfg,user=test")
1655
1656 // Test with invalid JSON
1657 sendDyncfgCmd(sd, "2-test",
1658 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "test"},
1659 []byte("{invalid json}"), "")
1660 },
1661 wantExposed: []wantExposedConfig{
1662 {
1663 discovererType: testDiscovererTypeNetListeners,
1664 name: "test-job",
1665 sourceType: "dyncfg",
1666 status: dyncfg.StatusAccepted,
1667 },
1668 },
1669 wantDyncfgFunc: func(t *testing.T, got string) {
1670 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
1671 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 2-test 400 application/json")
1672 assert.Contains(t, got, "Failed to parse config")
1673 },
1674 }
1675 },
1676 },
1677 }
1678
1679 for name, tc := range tests {
1680 t.Run(name, func(t *testing.T) {
1681 sim := tc.createSim()
1682 sim.run(t)
1683 })
1684 }
1685 }
1686
1687 func TestServiceDiscovery_DyncfgMultipleJobs(t *testing.T) {
1688 tests := map[string]struct {
1689 createSim func() *dyncfgSim
1690 }{
1691 "multiple jobs lifecycle": {
1692 createSim: func() *dyncfgSim {
1693 cfg1 := newTestNetListenersConfig("job1", 0, 0, defaultTestServices())
1694 cfg2 := newTestNetListenersConfig("job2", 0, 0, defaultTestServices())
1695 payload1, _ := json.Marshal(cfg1)
1696 payload2, _ := json.Marshal(cfg2)
1697
1698 return &dyncfgSim{
1699 do: func(sd *ServiceDiscovery) {
1700 // Add job1
1701 sendDyncfgCmd(sd, "1-add",
1702 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "job1"},
1703 payload1, "type=dyncfg,user=test")
1704
1705 // Add job2
1706 sendDyncfgCmd(sd, "2-add",
1707 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "job2"},
1708 payload2, "type=dyncfg,user=test")
1709
1710 // Enable both
1711 sendDyncfgCmd(sd, "3-enable",
1712 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "job1"), "enable"},
1713 nil, "")
1714 sendDyncfgCmd(sd, "4-enable",
1715 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "job2"), "enable"},
1716 nil, "")
1717
1718 // Disable job1
1719 sendDyncfgCmd(sd, "5-disable",
1720 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "job1"), "disable"},
1721 nil, "")
1722
1723 // Remove job2
1724 sendDyncfgCmd(sd, "6-remove",
1725 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "job2"), "remove"},
1726 nil, "")
1727 },
1728 wantExposed: []wantExposedConfig{
1729 {
1730 discovererType: testDiscovererTypeNetListeners,
1731 name: "job1",
1732 sourceType: "dyncfg",
1733 status: dyncfg.StatusDisabled,
1734 },
1735 },
1736 wantRunning: []string{},
1737 wantDyncfg: `
1738 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1739 {"status":202,"message":""}
1740 FUNCTION_RESULT_END
1741
1742 CONFIG test:sd:net_listeners:job1 create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1743
1744 FUNCTION_RESULT_BEGIN 2-add 202 application/json
1745 {"status":202,"message":""}
1746 FUNCTION_RESULT_END
1747
1748 CONFIG test:sd:net_listeners:job2 create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1749
1750 FUNCTION_RESULT_BEGIN 3-enable 200 application/json
1751 {"status":200,"message":""}
1752 FUNCTION_RESULT_END
1753
1754 CONFIG test:sd:net_listeners:job1 status running
1755
1756 FUNCTION_RESULT_BEGIN 4-enable 200 application/json
1757 {"status":200,"message":""}
1758 FUNCTION_RESULT_END
1759
1760 CONFIG test:sd:net_listeners:job2 status running
1761
1762 FUNCTION_RESULT_BEGIN 5-disable 200 application/json
1763 {"status":200,"message":""}
1764 FUNCTION_RESULT_END
1765
1766 CONFIG test:sd:net_listeners:job1 status disabled
1767
1768 FUNCTION_RESULT_BEGIN 6-remove 200 application/json
1769 {"status":200,"message":""}
1770 FUNCTION_RESULT_END
1771
1772 CONFIG test:sd:net_listeners:job2 delete
1773 `,
1774 }
1775 },
1776 },
1777 }
1778
1779 for name, tc := range tests {
1780 t.Run(name, func(t *testing.T) {
1781 sim := tc.createSim()
1782 sim.run(t)
1783 })
1784 }
1785 }
1786
1787 func TestServiceDiscovery_DyncfgPriority(t *testing.T) {
1788 tests := map[string]struct {
1789 createSim func() *dyncfgSim
1790 }{
1791 "dyncfg add replaces running file config": {
1792 // File config is running, dyncfg add with same name should replace it
1793 createSim: func() *dyncfgSim {
1794 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
1795 payload, _ := json.Marshal(cfg)
1796
1797 return &dyncfgSim{
1798 do: func(sd *ServiceDiscovery) {
1799 // Simulate a running file config
1800 fileCfg := sdConfig{
1801 "name": "test-job",
1802 ikeyDiscovererType: testDiscovererTypeNetListeners,
1803 ikeyPipelineKey: "/etc/netdata/sd.d/test.conf",
1804 ikeySource: "/etc/netdata/sd.d/test.conf",
1805 ikeySourceType: confgroup.TypeUser,
1806 }
1807 sd.seen.Add(fileCfg)
1808 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: fileCfg, Status: dyncfg.StatusRunning})
1809
1810 // Start the pipeline to simulate running state
1811 pipelineCfg := pipeline.Config{Name: "test-job"}
1812 _ = sd.mgr.Start(sd.ctx, fileCfg.PipelineKey(), pipelineCfg)
1813
1814 // Dyncfg add with same name - should replace file config
1815 sendDyncfgCmd(sd, "1-add",
1816 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1817 payload, "type=dyncfg,user=test")
1818 },
1819 wantExposed: []wantExposedConfig{
1820 {
1821 discovererType: testDiscovererTypeNetListeners,
1822 name: "test-job",
1823 sourceType: confgroup.TypeDyncfg, // dyncfg replaces file
1824 status: dyncfg.StatusAccepted,
1825 },
1826 },
1827 wantRunning: []string{}, // old pipeline stopped
1828 wantDyncfgFunc: func(t *testing.T, got string) {
1829 // Should see: create new dyncfg config (no delete - CONFIG create updates existing)
1830 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create accepted job")
1831 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
1832 },
1833 }
1834 },
1835 },
1836 "dyncfg add replaces stock file config": {
1837 // Stock file config exists, dyncfg add should replace it
1838 createSim: func() *dyncfgSim {
1839 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
1840 payload, _ := json.Marshal(cfg)
1841
1842 return &dyncfgSim{
1843 do: func(sd *ServiceDiscovery) {
1844 // Simulate a stock file config (not running)
1845 fileCfg := sdConfig{
1846 "name": "test-job",
1847 ikeyDiscovererType: testDiscovererTypeNetListeners,
1848 ikeyPipelineKey: "/usr/lib/netdata/conf.d/sd/test.conf",
1849 ikeySource: "/usr/lib/netdata/conf.d/sd/test.conf",
1850 ikeySourceType: confgroup.TypeStock,
1851 }
1852 sd.seen.Add(fileCfg)
1853 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: fileCfg, Status: dyncfg.StatusAccepted})
1854
1855 // Dyncfg add with same name - should replace stock config
1856 sendDyncfgCmd(sd, "1-add",
1857 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1858 payload, "type=dyncfg,user=test")
1859 },
1860 wantExposed: []wantExposedConfig{
1861 {
1862 discovererType: testDiscovererTypeNetListeners,
1863 name: "test-job",
1864 sourceType: confgroup.TypeDyncfg,
1865 status: dyncfg.StatusAccepted,
1866 },
1867 },
1868 wantDyncfgFunc: func(t *testing.T, got string) {
1869 // Should see: create new dyncfg config (no delete - CONFIG create updates existing)
1870 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create accepted job")
1871 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
1872 assert.Contains(t, got, "dyncfg")
1873 },
1874 }
1875 },
1876 },
1877 "dyncfg add replaces existing dyncfg config": {
1878 // Dyncfg config exists, another dyncfg add with same name should replace it (matching jobmgr pattern)
1879 createSim: func() *dyncfgSim {
1880 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
1881 payload, _ := json.Marshal(cfg)
1882
1883 return &dyncfgSim{
1884 do: func(sd *ServiceDiscovery) {
1885 // Simulate existing dyncfg config
1886 dyncfgCfg := sdConfig{
1887 "name": "test-job",
1888 ikeyDiscovererType: testDiscovererTypeNetListeners,
1889 ikeyPipelineKey: "dyncfg:net_listeners:test-job",
1890 ikeySource: "type=dyncfg,user=admin",
1891 ikeySourceType: confgroup.TypeDyncfg,
1892 }
1893 sd.seen.Add(dyncfgCfg)
1894 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: dyncfgCfg, Status: dyncfg.StatusRunning})
1895
1896 // Start the pipeline to simulate running state
1897 pipelineCfg := pipeline.Config{Name: "test-job"}
1898 _ = sd.mgr.Start(sd.ctx, dyncfgCfg.PipelineKey(), pipelineCfg)
1899
1900 // Another dyncfg add with same name - should replace (matching jobmgr pattern)
1901 sendDyncfgCmd(sd, "1-add",
1902 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1903 payload, "type=dyncfg,user=test")
1904 },
1905 wantExposed: []wantExposedConfig{
1906 {
1907 discovererType: testDiscovererTypeNetListeners,
1908 name: "test-job",
1909 sourceType: confgroup.TypeDyncfg,
1910 status: dyncfg.StatusAccepted, // new config in accepted state
1911 },
1912 },
1913 wantRunning: []string{}, // old pipeline stopped
1914 wantDyncfgFunc: func(t *testing.T, got string) {
1915 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-add 202 application/json")
1916 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create accepted job")
1917 },
1918 }
1919 },
1920 },
1921 }
1922
1923 for name, tc := range tests {
1924 t.Run(name, func(t *testing.T) {
1925 sim := tc.createSim()
1926 sim.run(t)
1927 })
1928 }
1929 }
1930
1931 func TestServiceDiscovery_DyncfgUpdateSameConfig(t *testing.T) {
1932 tests := map[string]struct {
1933 createSim func() *dyncfgSim
1934 }{
1935 "update running pipeline with same config skips restart": {
1936 createSim: func() *dyncfgSim {
1937 cfg := newTestNetListenersConfig("test-job", confopt.LongDuration(5*time.Second), 0, defaultTestServices())
1938 payload, _ := json.Marshal(cfg)
1939
1940 return &dyncfgSim{
1941 do: func(sd *ServiceDiscovery) {
1942 // Add
1943 sendDyncfgCmd(sd, "1-add",
1944 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
1945 payload, "type=dyncfg,user=test")
1946
1947 // Enable
1948 sendDyncfgCmd(sd, "2-enable",
1949 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
1950 nil, "")
1951
1952 // Update with exact same config (should return 200 without restart)
1953 sendDyncfgCmd(sd, "3-update",
1954 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
1955 payload, "type=dyncfg,user=test")
1956 },
1957 wantExposed: []wantExposedConfig{
1958 {
1959 discovererType: testDiscovererTypeNetListeners,
1960 name: "test-job",
1961 sourceType: "dyncfg",
1962 status: dyncfg.StatusRunning,
1963 },
1964 },
1965 wantRunning: []string{"dyncfg:net_listeners:test-job"},
1966 wantDyncfg: `
1967 FUNCTION_RESULT_BEGIN 1-add 202 application/json
1968 {"status":202,"message":""}
1969 FUNCTION_RESULT_END
1970
1971 CONFIG test:sd:net_listeners:test-job create accepted job /collectors/test/ServiceDiscovery dyncfg 'type=dyncfg,user=test' 'schema get enable disable update test userconfig remove' 0x0000 0x0000
1972
1973 FUNCTION_RESULT_BEGIN 2-enable 200 application/json
1974 {"status":200,"message":""}
1975 FUNCTION_RESULT_END
1976
1977 CONFIG test:sd:net_listeners:test-job status running
1978
1979 FUNCTION_RESULT_BEGIN 3-update 200 application/json
1980 {"status":200,"message":""}
1981 FUNCTION_RESULT_END
1982
1983 CONFIG test:sd:net_listeners:test-job status running
1984 `,
1985 }
1986 },
1987 },
1988 }
1989
1990 for name, tc := range tests {
1991 t.Run(name, func(t *testing.T) {
1992 sim := tc.createSim()
1993 sim.run(t)
1994 })
1995 }
1996 }
1997
1998 func TestServiceDiscovery_DyncfgUpdateFailedState(t *testing.T) {
1999 tests := map[string]struct {
2000 createSim func() *dyncfgSim
2001 }{
2002 "update config in failed state restarts pipeline": {
2003 createSim: func() *dyncfgSim {
2004 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
2005 updatedPayload, _ := json.Marshal(updatedCfg)
2006
2007 return &dyncfgSim{
2008 do: func(sd *ServiceDiscovery) {
2009 // Manually add a failed config
2010 failedCfg := sdConfig{
2011 "name": "test-job",
2012 ikeyDiscovererType: testDiscovererTypeNetListeners,
2013 ikeyPipelineKey: "dyncfg:net_listeners:test-job",
2014 ikeySource: "type=dyncfg,user=test",
2015 ikeySourceType: confgroup.TypeDyncfg,
2016 }
2017 sd.seen.Add(failedCfg)
2018 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: failedCfg, Status: dyncfg.StatusFailed})
2019
2020 // Update should restart the pipeline
2021 sendDyncfgCmd(sd, "1-update",
2022 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
2023 updatedPayload, "type=dyncfg,user=test")
2024 },
2025 wantExposed: []wantExposedConfig{
2026 {
2027 discovererType: testDiscovererTypeNetListeners,
2028 name: "test-job",
2029 sourceType: "dyncfg",
2030 status: dyncfg.StatusRunning,
2031 },
2032 },
2033 wantRunning: []string{"dyncfg:net_listeners:test-job"},
2034 wantDyncfgFunc: func(t *testing.T, got string) {
2035 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-update 200 application/json")
2036 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job status running")
2037 },
2038 }
2039 },
2040 },
2041 }
2042
2043 for name, tc := range tests {
2044 t.Run(name, func(t *testing.T) {
2045 sim := tc.createSim()
2046 sim.run(t)
2047 })
2048 }
2049 }
2050
2051 func TestServiceDiscovery_DyncfgEnableFromFailed(t *testing.T) {
2052 tests := map[string]struct {
2053 createSim func() *dyncfgSim
2054 }{
2055 "enable config from failed state starts pipeline": {
2056 createSim: func() *dyncfgSim {
2057 return &dyncfgSim{
2058 do: func(sd *ServiceDiscovery) {
2059 // Manually add a failed config with valid pipeline config data
2060 failedCfg := sdConfig{
2061 "name": "test-job",
2062 ikeyDiscovererType: testDiscovererTypeNetListeners,
2063 ikeyPipelineKey: "dyncfg:net_listeners:test-job",
2064 ikeySource: "type=dyncfg,user=test",
2065 ikeySourceType: confgroup.TypeDyncfg,
2066 "discoverer": map[string]any{
2067 "net_listeners": map[string]any{},
2068 },
2069 "services": []any{
2070 map[string]any{"id": "test-rule", "match": "true"},
2071 },
2072 }
2073 sd.seen.Add(failedCfg)
2074 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: failedCfg, Status: dyncfg.StatusFailed})
2075
2076 // Enable should start the pipeline
2077 sendDyncfgCmd(sd, "1-enable",
2078 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
2079 nil, "")
2080 },
2081 wantExposed: []wantExposedConfig{
2082 {
2083 discovererType: testDiscovererTypeNetListeners,
2084 name: "test-job",
2085 sourceType: "dyncfg",
2086 status: dyncfg.StatusRunning,
2087 },
2088 },
2089 wantRunning: []string{"dyncfg:net_listeners:test-job"},
2090 wantDyncfgFunc: func(t *testing.T, got string) {
2091 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-enable 200 application/json")
2092 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job status running")
2093 },
2094 }
2095 },
2096 },
2097 }
2098
2099 for name, tc := range tests {
2100 t.Run(name, func(t *testing.T) {
2101 sim := tc.createSim()
2102 sim.run(t)
2103 })
2104 }
2105 }
2106
2107 func TestServiceDiscovery_DyncfgConversionUpdate(t *testing.T) {
2108 tests := map[string]struct {
2109 createSim func() *dyncfgSim
2110 }{
2111 "update file config converts to dyncfg": {
2112 createSim: func() *dyncfgSim {
2113 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
2114 updatedPayload, _ := json.Marshal(updatedCfg)
2115
2116 return &dyncfgSim{
2117 do: func(sd *ServiceDiscovery) {
2118 // Simulate running file config
2119 fileCfg := sdConfig{
2120 "name": "test-job",
2121 ikeyDiscovererType: testDiscovererTypeNetListeners,
2122 ikeyPipelineKey: "/etc/netdata/sd.d/test.conf",
2123 ikeySource: "/etc/netdata/sd.d/test.conf",
2124 ikeySourceType: confgroup.TypeUser,
2125 "discoverer": map[string]any{
2126 "net_listeners": map[string]any{},
2127 },
2128 "services": []any{
2129 map[string]any{"id": "test-rule", "match": "true"},
2130 },
2131 }
2132 sd.seen.Add(fileCfg)
2133 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: fileCfg, Status: dyncfg.StatusRunning})
2134
2135 // Start the file pipeline
2136 pipelineCfg := pipeline.Config{Name: "test-job"}
2137 _ = sd.mgr.Start(sd.ctx, fileCfg.PipelineKey(), pipelineCfg)
2138
2139 // Update via dyncfg - should convert to dyncfg source
2140 sendDyncfgCmd(sd, "1-update",
2141 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
2142 updatedPayload, "type=dyncfg,user=admin")
2143 },
2144 wantExposed: []wantExposedConfig{
2145 {
2146 discovererType: testDiscovererTypeNetListeners,
2147 name: "test-job",
2148 sourceType: confgroup.TypeDyncfg, // Converted to dyncfg!
2149 status: dyncfg.StatusRunning,
2150 },
2151 },
2152 wantRunning: []string{"dyncfg:net_listeners:test-job"}, // New pipeline key
2153 wantDyncfgFunc: func(t *testing.T, got string) {
2154 // ConfigCreate acts as upsert (no delete needed)
2155 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create running job")
2156 assert.Contains(t, got, "dyncfg") // New source type
2157 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-update 200 application/json")
2158 },
2159 }
2160 },
2161 },
2162 "update disabled file config converts to dyncfg without starting": {
2163 createSim: func() *dyncfgSim {
2164 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
2165 updatedPayload, _ := json.Marshal(updatedCfg)
2166
2167 return &dyncfgSim{
2168 do: func(sd *ServiceDiscovery) {
2169 // Simulate disabled file config
2170 fileCfg := sdConfig{
2171 "name": "test-job",
2172 ikeyDiscovererType: testDiscovererTypeNetListeners,
2173 ikeyPipelineKey: "/etc/netdata/sd.d/test.conf",
2174 ikeySource: "/etc/netdata/sd.d/test.conf",
2175 ikeySourceType: confgroup.TypeUser,
2176 "discoverer": map[string]any{
2177 "net_listeners": map[string]any{},
2178 },
2179 "services": []any{
2180 map[string]any{"id": "test-rule", "match": "true"},
2181 },
2182 }
2183 sd.seen.Add(fileCfg)
2184 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: fileCfg, Status: dyncfg.StatusDisabled})
2185
2186 // Update via dyncfg - should convert but stay disabled
2187 sendDyncfgCmd(sd, "1-update",
2188 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
2189 updatedPayload, "type=dyncfg,user=admin")
2190 },
2191 wantExposed: []wantExposedConfig{
2192 {
2193 discovererType: testDiscovererTypeNetListeners,
2194 name: "test-job",
2195 sourceType: confgroup.TypeDyncfg,
2196 status: dyncfg.StatusDisabled, // Stays disabled
2197 },
2198 },
2199 wantRunning: []string{}, // Not running
2200 wantDyncfgFunc: func(t *testing.T, got string) {
2201 // ConfigCreate acts as upsert (no delete needed)
2202 assert.Contains(t, got, "CONFIG test:sd:net_listeners:test-job create disabled job")
2203 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 1-update 200 application/json")
2204 },
2205 }
2206 },
2207 },
2208 }
2209
2210 for name, tc := range tests {
2211 t.Run(name, func(t *testing.T) {
2212 sim := tc.createSim()
2213 sim.run(t)
2214 })
2215 }
2216 }
2217
2218 func TestServiceDiscovery_DyncfgRestartErrorHandling(t *testing.T) {
2219 tests := map[string]struct {
2220 createSim func() *dyncfgSim
2221 }{
2222 "restart with invalid config keeps old pipeline running": {
2223 createSim: func() *dyncfgSim {
2224 cfg := newTestNetListenersConfig("test-job", 0, 0, defaultTestServices())
2225 payload, _ := json.Marshal(cfg)
2226
2227 return &dyncfgSim{
2228 do: func(sd *ServiceDiscovery) {
2229 // Override newPipeline to fail on second call
2230 callCount := 0
2231 sd.newPipeline = func(cfg pipeline.Config) (sdPipeline, error) {
2232 callCount++
2233 if callCount > 1 {
2234 return nil, errors.New("simulated pipeline creation failure")
2235 }
2236 return newTestPipeline(cfg.Name), nil
2237 }
2238 // Also update mgr's newPipeline
2239 sd.mgr.newPipeline = sd.newPipeline
2240
2241 // Add and enable
2242 sendDyncfgCmd(sd, "1-add",
2243 []string{sd.dyncfgTemplateID(testDiscovererTypeNetListeners), "add", "test-job"},
2244 payload, "type=dyncfg,user=test")
2245
2246 sendDyncfgCmd(sd, "2-enable",
2247 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "enable"},
2248 nil, "")
2249
2250 // Update - pipeline creation will fail
2251 updatedCfg := newTestNetListenersConfig("test-job", confopt.LongDuration(10*time.Second), 0, defaultTestServices())
2252 updatedPayload, _ := json.Marshal(updatedCfg)
2253
2254 sendDyncfgCmd(sd, "3-update",
2255 []string{sd.dyncfgJobID(testDiscovererTypeNetListeners, "test-job"), "update"},
2256 updatedPayload, "type=dyncfg,user=test")
2257 },
2258 wantExposed: []wantExposedConfig{
2259 {
2260 discovererType: testDiscovererTypeNetListeners,
2261 name: "test-job",
2262 sourceType: "dyncfg",
2263 status: dyncfg.StatusRunning,
2264 },
2265 },
2266 // NOTE: When Restart fails preflight validation (newPipeline fails), the
2267 // old pipeline stays running and handler rolls back to old running state.
2268 wantRunning: []string{"dyncfg:net_listeners:test-job"},
2269 wantDyncfgFunc: func(t *testing.T, got string) {
2270 assert.Contains(t, got, "FUNCTION_RESULT_BEGIN 3-update 200 application/json")
2271 line := "CONFIG test:sd:net_listeners:test-job status running"
2272 assert.GreaterOrEqual(t, strings.Count(got, line), 2, "expected running status to be re-notified after update failure")
2273 assert.NotContains(t, got, "CONFIG test:sd:net_listeners:test-job status failed")
2274 },
2275 }
2276 },
2277 },
2278 }
2279
2280 for name, tc := range tests {
2281 t.Run(name, func(t *testing.T) {
2282 sim := tc.createSim()
2283 sim.run(t)
2284 })
2285 }
2286 }
2287
2288 func TestServiceDiscovery_DyncfgFileRemovalWithDyncfgOverride(t *testing.T) {
2289 tests := map[string]struct {
2290 createSim func() *dyncfgSim
2291 }{
2292 "file config removal does not affect dyncfg override": {
2293 createSim: func() *dyncfgSim {
2294 return &dyncfgSim{
2295 do: func(sd *ServiceDiscovery) {
2296 // Add file config to seen cache (simulating it was seen from file)
2297 fileCfg := sdConfig{
2298 "name": "test-job",
2299 ikeyDiscovererType: testDiscovererTypeNetListeners,
2300 ikeyPipelineKey: "/etc/netdata/sd.d/test.conf",
2301 ikeySource: "/etc/netdata/sd.d/test.conf",
2302 ikeySourceType: confgroup.TypeUser,
2303 }
2304 sd.seen.Add(fileCfg)
2305
2306 // Add dyncfg override (higher priority) to both caches
2307 dyncfgCfg := sdConfig{
2308 "name": "test-job",
2309 ikeyDiscovererType: testDiscovererTypeNetListeners,
2310 ikeyPipelineKey: "dyncfg:net_listeners:test-job",
2311 ikeySource: "type=dyncfg,user=test",
2312 ikeySourceType: confgroup.TypeDyncfg,
2313 }
2314 sd.seen.Add(dyncfgCfg)
2315 sd.exposed.Add(&dyncfg.Entry[sdConfig]{Cfg: dyncfgCfg, Status: dyncfg.StatusRunning})
2316
2317 // Start the dyncfg pipeline
2318 pipelineCfg := pipeline.Config{Name: "test-job"}
2319 _ = sd.mgr.Start(sd.ctx, dyncfgCfg.PipelineKey(), pipelineCfg)
2320
2321 // Simulate file removal by calling removePipeline
2322 sd.removePipeline(confFile{source: "/etc/netdata/sd.d/test.conf"})
2323 },
2324 wantExposed: []wantExposedConfig{
2325 {
2326 discovererType: testDiscovererTypeNetListeners,
2327 name: "test-job",
2328 sourceType: confgroup.TypeDyncfg, // Dyncfg still exposed
2329 status: dyncfg.StatusRunning,
2330 },
2331 },
2332 wantRunning: []string{"dyncfg:net_listeners:test-job"}, // Dyncfg pipeline still running
2333 }
2334 },
2335 },
2336 }
2337
2338 for name, tc := range tests {
2339 t.Run(name, func(t *testing.T) {
2340 sim := tc.createSim()
2341 sim.run(t)
2342 })
2343 }
2344 }
2345
2346 // testPipeline is a simple pipeline for testing that just waits for cancellation.
2347 type testPipeline struct {
2348 name string
2349 }
2350
2351 func newTestPipeline(name string) *testPipeline {
2352 return &testPipeline{name: name}
2353 }
2354
2355 func (p *testPipeline) Run(ctx context.Context, out chan<- []*confgroup.Group) {
2356 <-ctx.Done()
2357 }