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