| 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 | } |