master
go 183 lines 3.38 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package discovery
4
5 import (
6 "context"
7 "errors"
8 "fmt"
9 "log/slog"
10 "sync"
11 "time"
12
13 "github.com/netdata/netdata/go/plugins/logger"
14 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
15 )
16
17 func NewManager(cfg Config) (*Manager, error) {
18 if err := validateConfig(cfg); err != nil {
19 return nil, fmt.Errorf("discovery manager config validation: %v", err)
20 }
21
22 mgr := &Manager{
23 Logger: logger.New().With(
24 slog.String("component", "discovery manager"),
25 ),
26 send: make(chan struct{}, 1),
27 sendEvery: time.Second * 2, // timeout to aggregate changes
28 discoverers: make([]Discoverer, 0),
29 mux: &sync.RWMutex{},
30 cache: newCache(),
31 }
32
33 if err := mgr.registerDiscoverers(cfg); err != nil {
34 return nil, fmt.Errorf("discovery manager initializaion: %v", err)
35 }
36
37 return mgr, nil
38 }
39
40 type Manager struct {
41 *logger.Logger
42 discoverers []Discoverer
43 send chan struct{}
44 sendEvery time.Duration
45 mux *sync.RWMutex
46 cache *cache
47 }
48
49 func (m *Manager) String() string {
50 return fmt.Sprintf("discovery manager: %v", m.discoverers)
51 }
52
53 func (m *Manager) Run(ctx context.Context, in chan<- []*confgroup.Group) {
54 m.Info("instance is started")
55 defer func() { m.Info("instance is stopped") }()
56
57 var wg sync.WaitGroup
58
59 for _, d := range m.discoverers {
60 wg.Add(1)
61 go func(d Discoverer) {
62 defer wg.Done()
63 m.runDiscoverer(ctx, d)
64 }(d)
65 }
66
67 wg.Go(func() {
68 m.sendLoop(ctx, in)
69 })
70
71 wg.Wait()
72 <-ctx.Done()
73 }
74
75 func (m *Manager) registerDiscoverers(cfg Config) error {
76 for _, provider := range cfg.Providers {
77 d, enabled, err := provider.Build(cfg.BuildContext)
78 if err != nil {
79 return fmt.Errorf("provider '%s': %w", provider.Name(), err)
80 }
81 if !enabled {
82 continue
83 }
84 if d == nil {
85 return fmt.Errorf("provider '%s': enabled but returned nil discoverer", provider.Name())
86 }
87 m.discoverers = append(m.discoverers, d)
88 }
89
90 if len(m.discoverers) == 0 {
91 return errors.New("zero registered discoverers")
92 }
93
94 m.Infof("registered discoverers: %v", m.discoverers)
95
96 return nil
97 }
98
99 func (m *Manager) runDiscoverer(ctx context.Context, d Discoverer) {
100 done := make(chan struct{})
101 updates := make(chan []*confgroup.Group)
102
103 go func() { defer close(done); d.Run(ctx, updates) }()
104
105 for {
106 select {
107 case <-ctx.Done():
108 select {
109 case <-done:
110 case <-time.After(time.Second * 10):
111 }
112 return
113 case groups, ok := <-updates:
114 if !ok {
115 return
116 }
117 func() {
118 m.mux.Lock()
119 defer m.mux.Unlock()
120
121 m.cache.update(groups)
122 m.triggerSend()
123 }()
124 }
125 }
126 }
127
128 func (m *Manager) sendLoop(ctx context.Context, in chan<- []*confgroup.Group) {
129 m.mustSend(ctx, in)
130
131 tk := time.NewTicker(m.sendEvery)
132 defer tk.Stop()
133
134 for {
135 select {
136 case <-ctx.Done():
137 return
138 case <-tk.C:
139 select {
140 case <-m.send:
141 m.trySend(in)
142 default:
143 }
144 }
145 }
146 }
147
148 func (m *Manager) mustSend(ctx context.Context, in chan<- []*confgroup.Group) {
149 select {
150 case <-ctx.Done():
151 return
152 case <-m.send:
153 m.mux.Lock()
154 groups := m.cache.groups()
155 m.cache.reset()
156 m.mux.Unlock()
157
158 select {
159 case <-ctx.Done():
160 case in <- groups:
161 }
162 return
163 }
164 }
165
166 func (m *Manager) trySend(in chan<- []*confgroup.Group) {
167 m.mux.Lock()
168 defer m.mux.Unlock()
169
170 select {
171 case in <- m.cache.groups():
172 m.cache.reset()
173 default:
174 m.triggerSend()
175 }
176 }
177
178 func (m *Manager) triggerSend() {
179 select {
180 case m.send <- struct{}{}:
181 default:
182 }
183 }