master
go 104 lines 2.04 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package file
4
5 import (
6 "context"
7 "errors"
8 "fmt"
9 "log/slog"
10 "sync"
11
12 "github.com/netdata/netdata/go/plugins/logger"
13 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
14 )
15
16 var log = logger.New().With(
17 slog.String("component", "discovery"),
18 slog.String("discoverer", "file"),
19 )
20
21 func NewDiscovery(cfg Config) (*Discovery, error) {
22 if err := validateConfig(cfg); err != nil {
23 return nil, fmt.Errorf("file discovery config validation: %v", err)
24 }
25
26 d := Discovery{
27 Logger: log,
28 }
29
30 if err := d.registerDiscoverers(cfg); err != nil {
31 return nil, fmt.Errorf("file discovery initialization: %v", err)
32 }
33
34 return &d, nil
35 }
36
37 type (
38 Discovery struct {
39 *logger.Logger
40 discoverers []discoverer
41 }
42 discoverer interface {
43 Run(ctx context.Context, in chan<- []*confgroup.Group)
44 }
45 )
46
47 func (d *Discovery) String() string {
48 return d.Name()
49 }
50
51 func (d *Discovery) Name() string {
52 return fmt.Sprintf("file discovery: %v", d.discoverers)
53 }
54
55 func (d *Discovery) registerDiscoverers(cfg Config) error {
56 if len(cfg.Read) != 0 {
57 d.discoverers = append(d.discoverers, NewReader(cfg.Registry, cfg.Read))
58 }
59 if len(cfg.Watch) != 0 {
60 d.discoverers = append(d.discoverers, NewWatcher(cfg.Registry, cfg.Watch))
61 }
62 if len(d.discoverers) == 0 {
63 return errors.New("zero registered discoverers")
64 }
65 return nil
66 }
67
68 func (d *Discovery) Run(ctx context.Context, in chan<- []*confgroup.Group) {
69 d.Info("instance is started")
70 defer func() { d.Info("instance is stopped") }()
71
72 var wg sync.WaitGroup
73
74 for _, dd := range d.discoverers {
75 wg.Add(1)
76 go func(dd discoverer) {
77 defer wg.Done()
78 d.runDiscoverer(ctx, dd, in)
79 }(dd)
80 }
81
82 wg.Wait()
83 <-ctx.Done()
84 }
85
86 func (d *Discovery) runDiscoverer(ctx context.Context, dd discoverer, in chan<- []*confgroup.Group) {
87 updates := make(chan []*confgroup.Group)
88 go dd.Run(ctx, updates)
89 for {
90 select {
91 case <-ctx.Done():
92 return
93 case groups, ok := <-updates:
94 if !ok {
95 return
96 }
97 select {
98 case <-ctx.Done():
99 return
100 case in <- groups:
101 }
102 }
103 }
104 }