master
go 376 lines 11 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package sd
4
5 import (
6 "context"
7 "fmt"
8 "strings"
9
10 "github.com/netdata/netdata/go/plugins/pkg/netdataapi"
11 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
12 "github.com/netdata/netdata/go/plugins/plugin/framework/dyncfg"
13 "github.com/netdata/netdata/go/plugins/plugin/framework/functions"
14 )
15
16 const (
17 dyncfgSDPrefixf = "%s:sd:"
18 dyncfgSDPath = "/collectors/%s/ServiceDiscovery"
19 )
20
21 func (d *ServiceDiscovery) dyncfgSDPrefixValue() string {
22 return fmt.Sprintf(dyncfgSDPrefixf, d.pluginName)
23 }
24
25 func (d *ServiceDiscovery) dyncfgTemplateID(discovererType string) string {
26 return fmt.Sprintf("%s%s", d.dyncfgSDPrefixValue(), discovererType)
27 }
28
29 func (d *ServiceDiscovery) dyncfgJobID(discovererType, name string) string {
30 return fmt.Sprintf("%s%s:%s", d.dyncfgSDPrefixValue(), discovererType, name)
31 }
32
33 func dyncfgSDTemplateCmds() string {
34 return dyncfg.JoinCommands(
35 dyncfg.CommandAdd,
36 dyncfg.CommandSchema,
37 dyncfg.CommandTest,
38 dyncfg.CommandUserconfig,
39 )
40 }
41
42 func dyncfgTemplateJobName(fn dyncfg.Function) string {
43 if name := fn.JobName(); name != "" {
44 return name
45 }
46 return "test"
47 }
48
49 func (d *ServiceDiscovery) dyncfgSDTemplateCreate(discovererType string) {
50 d.dyncfgApi.ConfigCreate(netdataapi.ConfigOpts{
51 ID: d.dyncfgTemplateID(discovererType),
52 Status: dyncfg.StatusAccepted.String(),
53 ConfigType: dyncfg.ConfigTypeTemplate.String(),
54 Path: fmt.Sprintf(dyncfgSDPath, d.pluginName),
55 SourceType: "internal",
56 Source: "internal",
57 SupportedCommands: dyncfgSDTemplateCmds(),
58 })
59 }
60
61 // sdCallbacks implements dyncfg.Callbacks[sdConfig]
62 type sdCallbacks struct {
63 sd *ServiceDiscovery
64 }
65
66 func (cb *sdCallbacks) ExtractKey(fn dyncfg.Function) (key, name string, ok bool) {
67 id := fn.ID()
68 if fn.Command() == dyncfg.CommandAdd {
69 dt, _, _ := cb.sd.extractDiscovererAndName(id)
70 if dt == "" || !cb.sd.hasDiscovererType(dt) {
71 return "", "", false
72 }
73 name = fn.JobName()
74 if name == "" {
75 return "", "", false
76 }
77 return dt + ":" + name, name, true
78 }
79 dt, name, isJob := cb.sd.extractDiscovererAndName(id)
80 if !isJob || name == "" {
81 return "", "", false
82 }
83 return dt + ":" + name, name, true
84 }
85
86 func (cb *sdCallbacks) ValidateJobName(name string) error {
87 return dyncfg.JobNameRuleAllowDots(name)
88 }
89
90 func (cb *sdCallbacks) ParseAndValidate(fn dyncfg.Function, name string) (sdConfig, error) {
91 dt, _, _ := cb.sd.extractDiscovererAndName(fn.ID())
92 if _, err := parseDyncfgPayload(fn.Payload(), dt, name, cb.sd.configDefaults, cb.sd.discovererRegistry(), true); err != nil {
93 return nil, err
94 }
95 pkey := pipelineKey(dt, name)
96 cfg, err := newSDConfigFromJSON(fn.Payload(), name, fn.Source(), confgroup.TypeDyncfg, dt, pkey)
97 if err != nil {
98 return nil, err
99 }
100 return cfg, nil
101 }
102
103 func (cb *sdCallbacks) Start(cfg sdConfig) error {
104 pipelineCfg, err := cfg.ToPipelineConfig(cb.sd.configDefaults)
105 if err != nil {
106 return err
107 }
108 return cb.sd.mgr.Start(cb.sd.ctx, cfg.PipelineKey(), pipelineCfg)
109 }
110
111 func (cb *sdCallbacks) Update(oldCfg, newCfg sdConfig) error {
112 pipelineCfg, err := newCfg.ToPipelineConfig(cb.sd.configDefaults)
113 if err != nil {
114 return dyncfg.MarkNonDisruptiveUpdate(err)
115 }
116 return cb.sd.mgr.Restart(cb.sd.ctx, newCfg.PipelineKey(), pipelineCfg)
117 }
118
119 func (cb *sdCallbacks) Stop(cfg sdConfig) {
120 cb.sd.mgr.Stop(cfg.PipelineKey())
121 }
122
123 func (cb *sdCallbacks) OnStatusChange(_ *dyncfg.Entry[sdConfig], _ dyncfg.Status, _ dyncfg.Function) {
124 }
125
126 func (cb *sdCallbacks) ConfigID(cfg sdConfig) string {
127 return cb.sd.dyncfgJobID(cfg.DiscovererType(), cfg.Name())
128 }
129
130 // dyncfgConfig is the handler for dyncfg config commands.
131 // Read-only commands (schema, get, userconfig) are executed directly.
132 // State-changing commands are queued for serial execution.
133 func (d *ServiceDiscovery) dyncfgConfig(fn dyncfg.Function) {
134 if err := fn.ValidateArgs(2); err != nil {
135 d.Warningf("dyncfg: %v", err)
136 d.dyncfgApi.SendCodef(fn, 400, "%v", err)
137 return
138 }
139
140 // Read-only commands can be executed directly
141 switch fn.Command() {
142 case dyncfg.CommandSchema:
143 d.dyncfgCmdSchema(fn)
144 return
145 case dyncfg.CommandGet:
146 d.dyncfgCmdGet(fn)
147 return
148 case dyncfg.CommandUserconfig:
149 d.dyncfgCmdUserconfig(fn)
150 return
151 case dyncfg.CommandTest:
152 // Test command validates config without creating a job
153 d.dyncfgCmdTest(fn)
154 return
155 }
156
157 // State-changing commands are queued for serial execution.
158 d.enqueueDyncfgFunction(fn)
159 }
160
161 // dyncfgSeqExec executes state-changing dyncfg commands serially.
162 func (d *ServiceDiscovery) dyncfgSeqExec(fn dyncfg.Function) {
163 d.handler.SyncDecision(fn)
164
165 switch fn.Command() {
166 case dyncfg.CommandAdd:
167 d.handler.CmdAdd(fn)
168 case dyncfg.CommandUpdate:
169 d.handler.CmdUpdate(fn)
170 case dyncfg.CommandEnable:
171 d.handler.CmdEnable(fn)
172 case dyncfg.CommandDisable:
173 d.handler.CmdDisable(fn)
174 case dyncfg.CommandRemove:
175 d.handler.CmdRemove(fn)
176 default:
177 d.Warningf("dyncfg: command '%s' not implemented", fn.Command())
178 d.dyncfgApi.SendCodef(fn, 501, "Command '%s' is not implemented.", fn.Command())
179 }
180 }
181
182 // dyncfgCmdSchema handles the schema command for templates and jobs
183 func (d *ServiceDiscovery) dyncfgCmdSchema(fn dyncfg.Function) {
184 id := fn.ID()
185 dt, _, _ := d.extractDiscovererAndName(id)
186
187 if dt == "" {
188 d.Warningf("dyncfg: schema: invalid ID format '%s'", id)
189 d.dyncfgApi.SendCodef(fn, 400, "Invalid ID format: %s", id)
190 return
191 }
192
193 desc, ok := d.discovererRegistry().Get(dt)
194 if !ok {
195 d.Warningf("dyncfg: schema: unknown discoverer type '%s'", dt)
196 d.dyncfgApi.SendCodef(fn, 404, "Unknown discoverer type: %s", dt)
197 return
198 }
199
200 schema := desc.Schema
201 if schema == "" {
202 d.Warningf("dyncfg: schema: discoverer type '%s' has no schema configured", dt)
203 d.dyncfgApi.SendCodef(fn, 500, "Schema is not configured for discoverer type: %s", dt)
204 return
205 }
206 d.dyncfgApi.SendJSON(fn, schema)
207 }
208
209 // dyncfgCmdGet handles the get command for jobs
210 func (d *ServiceDiscovery) dyncfgCmdGet(fn dyncfg.Function) {
211 id := fn.ID()
212 dt, name, isJob := d.extractDiscovererAndName(id)
213
214 if !isJob || name == "" {
215 d.Warningf("dyncfg: get: invalid job ID format '%s'", id)
216 d.dyncfgApi.SendCodef(fn, 400, "Invalid job ID format: %s", id)
217 return
218 }
219
220 entry, ok := d.exposed.LookupByKey(dt + ":" + name)
221 if !ok {
222 d.Warningf("dyncfg: get: config '%s:%s' not found", dt, name)
223 d.dyncfgApi.SendCodef(fn, 404, "Config '%s:%s' not found.", dt, name)
224 return
225 }
226
227 // Convert stored config to JSON via typed struct for consistent field ordering
228 bs, err := configToJSON(entry.Cfg.DataJSON())
229 if err != nil {
230 d.Warningf("dyncfg: get: failed to convert config '%s:%s' to JSON: %v", dt, name, err)
231 d.dyncfgApi.SendCodef(fn, 500, "Failed to convert config to JSON: %v", err)
232 return
233 }
234
235 d.dyncfgApi.SendJSON(fn, string(bs))
236 }
237
238 // dyncfgCmdTest handles the test command for templates and jobs (validates config without applying it)
239 func (d *ServiceDiscovery) dyncfgCmdTest(fn dyncfg.Function) {
240 id := fn.ID()
241
242 dt, name, isJob := d.extractDiscovererAndName(id)
243 if dt == "" || !d.hasDiscovererType(dt) {
244 d.Warningf("dyncfg: test: invalid discoverer type in ID '%s'", id)
245 d.dyncfgApi.SendCodef(fn, 400, "Invalid discoverer type in ID: %s", id)
246 return
247 }
248
249 if err := fn.ValidateHasPayload(); err != nil {
250 d.Warningf("dyncfg: test: %v for '%s'", err, dt)
251 d.dyncfgApi.SendCodef(fn, 400, "%v", err)
252 return
253 }
254
255 if !isJob {
256 name = dyncfgTemplateJobName(fn)
257 }
258 if err := dyncfg.JobNameRuleAllowDots(name); err != nil {
259 d.Warningf("dyncfg: test: unacceptable job name '%s' for '%s': %v", name, dt, err)
260 d.dyncfgApi.SendCodef(fn, 400, "Unacceptable job name '%s': %v.", name, err)
261 return
262 }
263
264 // Parse and validate the config without storing it
265 _, err := parseDyncfgPayload(fn.Payload(), dt, name, d.configDefaults, d.discovererRegistry(), true)
266 if err != nil {
267 d.Warningf("dyncfg: test: failed to parse config for '%s': %v", dt, err)
268 d.dyncfgApi.SendCodef(fn, 400, "Failed to parse config: %v", err)
269 return
270 }
271
272 if isJob {
273 d.Infof("dyncfg: test: config for '%s:%s' is valid", dt, name)
274 } else {
275 d.Infof("dyncfg: test: config for '%s' is valid", dt)
276 }
277 d.dyncfgApi.SendCodef(fn, 200, "")
278 }
279
280 // dyncfgCmdUserconfig handles the userconfig command for templates and jobs
281 // Returns YAML representation of the config for user-friendly file format
282 func (d *ServiceDiscovery) dyncfgCmdUserconfig(fn dyncfg.Function) {
283 id := fn.ID()
284 dt, name, isJob := d.extractDiscovererAndName(id)
285
286 if !d.hasDiscovererType(dt) {
287 d.Warningf("dyncfg: userconfig: invalid discoverer type in ID '%s'", id)
288 d.dyncfgApi.SendCodef(fn, 400, "Invalid discoverer type in ID: %s", id)
289 return
290 }
291
292 if !fn.HasPayload() {
293 d.Warningf("dyncfg: userconfig: missing payload for '%s'", id)
294 d.dyncfgApi.SendCodef(fn, 400, "Missing configuration payload.")
295 return
296 }
297
298 jobName := name
299 if !isJob || jobName == "" {
300 jobName = dyncfgTemplateJobName(fn)
301 }
302
303 if _, err := parseDyncfgPayload(fn.Payload(), dt, jobName, d.configDefaults, d.discovererRegistry(), false); err != nil {
304 d.Warningf("dyncfg: userconfig: failed to parse config for '%s': %v", id, err)
305 d.dyncfgApi.SendCodef(fn, 400, "Failed to parse config: %v", err)
306 return
307 }
308
309 bs, err := userConfigFromPayload(fn.Payload(), dt, jobName)
310 if err != nil {
311 d.Warningf("dyncfg: userconfig: failed to create config for '%s': %v", id, err)
312 d.dyncfgApi.SendCodef(fn, 400, "Failed to create config: %v", err)
313 return
314 }
315
316 d.dyncfgApi.SendYAML(fn, string(bs))
317 }
318
319 // extractDiscovererAndName parses a dyncfg ID into discoverer type and name.
320 // ID format: {prefix}{discovererType} (template) or {prefix}{discovererType}:{name} (job)
321 // Returns discovererType, name, isJob
322 func (d *ServiceDiscovery) extractDiscovererAndName(id string) (discovererType, name string, isJob bool) {
323 prefix := d.dyncfgSDPrefixValue()
324 if !strings.HasPrefix(id, prefix) {
325 return "", "", false
326 }
327
328 rest := strings.TrimPrefix(id, prefix)
329 if rest == "" {
330 return "", "", false
331 }
332
333 parts := strings.SplitN(rest, ":", 2)
334 discovererType = parts[0]
335
336 if len(parts) == 2 {
337 name = parts[1]
338 isJob = true
339 }
340
341 return discovererType, name, isJob
342 }
343
344 // registerDyncfgTemplates registers dyncfg templates for each discoverer type
345 func (d *ServiceDiscovery) registerDyncfgTemplates(ctx context.Context) {
346 if d.fnReg == nil {
347 return
348 }
349
350 // Register prefix handler for config commands
351 // Wrap to convert functions.Function to dyncfg.Function
352 d.fnReg.RegisterPrefix("config", d.dyncfgSDPrefixValue(), dyncfg.WrapHandler(d.dyncfgConfig))
353
354 // Register templates for each discoverer type
355 for _, dt := range d.discovererRegistry().Types() {
356 d.dyncfgSDTemplateCreate(dt)
357 d.Infof("registered dyncfg template for discoverer type '%s'", dt)
358 }
359 }
360
361 // unregisterDyncfgTemplates unregisters dyncfg templates
362 func (d *ServiceDiscovery) unregisterDyncfgTemplates() {
363 if d.fnReg == nil {
364 return
365 }
366
367 d.fnReg.UnregisterPrefix("config", d.dyncfgSDPrefixValue())
368 }
369
370 // autoEnableConfig enables a config without waiting for netdata's enable command.
371 func (d *ServiceDiscovery) autoEnableConfig(cfg sdConfig) {
372 fn := dyncfg.NewFunction(functions.Function{
373 Args: []string{d.dyncfgJobID(cfg.DiscovererType(), cfg.Name()), "enable"},
374 })
375 d.handler.CmdEnable(fn)
376 }