master
go 203 lines 4.51 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package pipeline
4
5 import (
6 "context"
7 "errors"
8 "fmt"
9 "log/slog"
10 "time"
11
12 "github.com/netdata/netdata/go/plugins/logger"
13 "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/model"
14 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
15 )
16
17 type DiscovererFactory func(payload DiscovererPayload, source string) ([]model.Discoverer, error)
18
19 func New(cfg Config, makeDiscoverers DiscovererFactory) (*Pipeline, error) {
20 if err := ValidateConfig(cfg); err != nil {
21 return nil, err
22 }
23 if makeDiscoverers == nil {
24 return nil, errors.New("discoverer factory is not set")
25 }
26
27 p := &Pipeline{
28 Logger: logger.New().With(
29 slog.String("component", "service discovery"),
30 slog.String("pipeline", cfg.Name),
31 ),
32 configDefaults: cfg.ConfigDefaults,
33 accum: newAccumulator(),
34 discoverers: make([]model.Discoverer, 0),
35 configs: make(map[string]map[uint64][]confgroup.Config),
36 }
37
38 p.accum.Logger = p.Logger
39
40 svr, err := newServiceEngine(cfg.Services)
41 if err != nil {
42 return nil, fmt.Errorf("services rules: %v", err)
43 }
44 p.svr = svr
45 svr.Logger = p.Logger
46
47 if err := p.registerDiscoverers(cfg, makeDiscoverers); err != nil {
48 return nil, err
49 }
50
51 return p, nil
52 }
53
54 type (
55 Pipeline struct {
56 *logger.Logger
57
58 configDefaults confgroup.Registry
59 discoverers []model.Discoverer
60 accum *accumulator
61
62 configs map[string]map[uint64][]confgroup.Config // [targetSource][targetHash]
63
64 // new
65 svr composer
66 }
67 composer interface {
68 compose(model.Target) []confgroup.Config
69 }
70 )
71
72 func (p *Pipeline) registerDiscoverers(conf Config, makeDiscoverers DiscovererFactory) error {
73 discs, err := makeDiscoverers(conf.Discoverer, conf.Source)
74 if err != nil {
75 return fmt.Errorf("failed to create %q discoverer: %w", conf.Discoverer.Type(), err)
76 }
77 p.discoverers = append(p.discoverers, discs...)
78
79 if len(p.discoverers) == 0 {
80 return errors.New("no discoverers registered")
81 }
82
83 return nil
84 }
85
86 func (p *Pipeline) Run(ctx context.Context, in chan<- []*confgroup.Group) {
87 p.Info("instance is started")
88 defer p.Info("instance is stopped")
89
90 p.accum.discoverers = p.discoverers
91
92 updates := make(chan []model.TargetGroup)
93 done := make(chan struct{})
94
95 go func() { defer close(done); p.accum.run(ctx, updates) }()
96
97 for {
98 select {
99 case <-ctx.Done():
100 select {
101 case <-done:
102 case <-time.After(time.Second * 10):
103 }
104 return
105 case <-done:
106 return
107 case tggs := <-updates:
108 p.Debugf("received %d target groups", len(tggs))
109 if cfggs := p.processGroups(tggs); len(cfggs) > 0 {
110 select {
111 case <-ctx.Done():
112 case in <- cfggs: // FIXME: potentially stale configs if upstream cannot receive (blocking)
113 }
114 }
115 }
116 }
117 }
118
119 func (p *Pipeline) processGroups(tggs []model.TargetGroup) []*confgroup.Group {
120 var groups []*confgroup.Group
121 // updates come from the accumulator; this ensures that all groups have different sources
122 for _, tgg := range tggs {
123 p.Debugf("processing group '%s' with %d target(s)", tgg.Source(), len(tgg.Targets()))
124 if v := p.processGroup(tgg); v != nil {
125 groups = append(groups, v)
126 }
127 }
128 return groups
129 }
130
131 func (p *Pipeline) processGroup(tgg model.TargetGroup) *confgroup.Group {
132 if len(tgg.Targets()) == 0 {
133 if _, ok := p.configs[tgg.Source()]; !ok {
134 return nil
135 }
136 delete(p.configs, tgg.Source())
137
138 return &confgroup.Group{Source: tgg.Source()}
139 }
140
141 targetsCache, ok := p.configs[tgg.Source()]
142 if !ok {
143 targetsCache = make(map[uint64][]confgroup.Config)
144 p.configs[tgg.Source()] = targetsCache
145 }
146
147 var changed bool
148 seen := make(map[uint64]bool)
149
150 for _, tgt := range tgg.Targets() {
151 if tgt == nil {
152 continue
153 }
154
155 hash := tgt.Hash()
156 seen[hash] = true
157
158 if _, ok := targetsCache[hash]; ok {
159 continue
160 }
161
162 targetsCache[hash] = nil
163
164 if p.svr != nil {
165 if cfgs := p.svr.compose(tgt); len(cfgs) > 0 {
166 targetsCache[hash] = cfgs
167 changed = true
168 for _, cfg := range cfgs {
169 cfg.SetProvider(tgg.Provider())
170 cfg.SetSource(tgg.Source())
171 cfg.SetSourceType(confgroup.TypeDiscovered)
172 if def, ok := p.configDefaults.Lookup(cfg.Module()); ok {
173 cfg.ApplyDefaults(def)
174 }
175 }
176 }
177 continue
178 }
179 }
180
181 for hash := range targetsCache {
182 if seen[hash] {
183 continue
184 }
185 if cfgs := targetsCache[hash]; len(cfgs) > 0 {
186 changed = true
187 }
188 delete(targetsCache, hash)
189 }
190
191 if !changed {
192 return nil
193 }
194
195 // TODO: deepcopy?
196 cfgGroup := &confgroup.Group{Source: tgg.Source()}
197
198 for _, cfgs := range targetsCache {
199 cfgGroup.Configs = append(cfgGroup.Configs, cfgs...)
200 }
201
202 return cfgGroup
203 }