master
go 525 lines 14.6 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package jobmgr
4
5 import (
6 "context"
7 "fmt"
8 "io"
9 "log/slog"
10 "slices"
11 "sync"
12 "time"
13
14 "github.com/netdata/netdata/go/plugins/logger"
15 "github.com/netdata/netdata/go/plugins/pkg/netdataapi"
16 "github.com/netdata/netdata/go/plugins/pkg/ticker"
17 "github.com/netdata/netdata/go/plugins/plugin/agent/jobmgr/funcctl"
18 "github.com/netdata/netdata/go/plugins/plugin/agent/jobmgr/secretsctl"
19 "github.com/netdata/netdata/go/plugins/plugin/agent/jobmgr/vnodectl"
20 "github.com/netdata/netdata/go/plugins/plugin/agent/policy"
21 "github.com/netdata/netdata/go/plugins/plugin/agent/secrets/resolver"
22 "github.com/netdata/netdata/go/plugins/plugin/agent/secrets/secretstore"
23 "github.com/netdata/netdata/go/plugins/plugin/agent/secrets/secretstore/backends"
24 "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi"
25 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
26 "github.com/netdata/netdata/go/plugins/plugin/framework/dyncfg"
27 "github.com/netdata/netdata/go/plugins/plugin/framework/functions"
28 "github.com/netdata/netdata/go/plugins/plugin/framework/metricsaudit"
29 "github.com/netdata/netdata/go/plugins/plugin/framework/runtimecomp"
30 "github.com/netdata/netdata/go/plugins/plugin/framework/vnoderegistry"
31 "github.com/netdata/netdata/go/plugins/plugin/framework/vnodes"
32 )
33
34 type Config struct {
35 PluginName string
36 Out io.Writer
37 RunModePolicy policy.RunModePolicy
38 Modules collectorapi.Registry
39 RunJob []string
40 ConfigDefaults confgroup.Registry
41 VarLibDir string
42 FnReg FunctionRegistry
43 Vnodes map[string]*vnodes.VirtualNode
44 SecretStores []secretstore.Config
45 SecretStoreService secretstore.Service
46 AuditMode bool
47 AuditAnalyzer metricsaudit.Analyzer
48 AuditDataDir string
49 FunctionJSONWriter func(payload []byte, code int)
50 RuntimeService runtimecomp.Service
51 VnodeRegistry *vnoderegistry.Registry
52 }
53
54 const (
55 cmdTestWorkerCap = 4
56 cmdTestDefaultTimeout = 60 * time.Second
57 cmdTestWorkerDrainWait = 5 * time.Second
58 )
59
60 func New(cfg Config) *Manager {
61 out := cfg.Out
62 if out == nil {
63 out = io.Discard
64 }
65
66 seen := dyncfg.NewSeenCache[confgroup.Config]()
67 exposed := dyncfg.NewExposedCache[confgroup.Config]()
68 api := dyncfg.NewResponder(netdataapi.New(out))
69 fnReg := cfg.FnReg
70 if fnReg == nil {
71 fnReg = noop{}
72 }
73 if provider, ok := fnReg.(interface {
74 TerminalFinalizer() functions.TerminalFinalizer
75 }); ok {
76 api.SetTerminalFinalizer(provider.TerminalFinalizer())
77 }
78 vnodesReg := cfg.Vnodes
79 if vnodesReg == nil {
80 vnodesReg = make(map[string]*vnodes.VirtualNode)
81 }
82 secretStoreSvc := cfg.SecretStoreService
83 if secretStoreSvc == nil {
84 storeCreators := backends.Creators()
85 secretStoreSvc = secretstore.NewService(storeCreators...)
86 }
87 vnodeRegistry := cfg.VnodeRegistry
88 if vnodeRegistry == nil {
89 vnodeRegistry = vnoderegistry.New()
90 }
91
92 mgr := &Manager{
93 Logger: logger.New().With(
94 slog.String("component", "job manager"),
95 ),
96 pluginName: cfg.PluginName,
97 out: out,
98 runModePolicy: cfg.RunModePolicy,
99 modules: cfg.Modules,
100 runJobNames: cfg.RunJob,
101 configDefaults: cfg.ConfigDefaults,
102 varLibDir: cfg.VarLibDir,
103 fnReg: fnReg,
104 initialSecretStores: append([]secretstore.Config(nil), cfg.SecretStores...),
105
106 auditMode: cfg.AuditMode,
107 auditAnalyzer: cfg.AuditAnalyzer,
108 auditDataDir: cfg.AuditDataDir,
109
110 discoveredConfigs: newDiscoveredConfigsCache(),
111 collectorSeen: seen,
112 collectorExposed: exposed,
113 secretStoreDeps: newSecretStoreDeps(),
114 runningJobs: newRunningJobsCache(),
115 retryingTasks: newRetryingTasksCache(),
116
117 started: make(chan struct{}),
118 addCh: make(chan confgroup.Config),
119 rmCh: make(chan confgroup.Config),
120 dyncfgCh: make(chan dyncfg.Function, 32),
121 cmdTestSem: make(chan struct{}, cmdTestWorkerCap),
122
123 dyncfgResponder: api,
124 runtimeService: cfg.RuntimeService,
125 vnodeRegistry: vnodeRegistry,
126 secretResolver: secretresolver.New(),
127 }
128 mgr.funcCtl = funcctl.New(funcctl.Options{
129 Logger: mgr.Logger,
130 FnReg: fnReg,
131 API: api,
132 JSONWriter: cfg.FunctionJSONWriter,
133 })
134
135 mgr.collectorCallbacks = &collectorCallbacks{mgr: mgr}
136 mgr.collectorHandler = dyncfg.NewHandler(dyncfg.HandlerOpts[confgroup.Config]{
137 Logger: mgr.Logger,
138 API: api,
139 Seen: seen,
140 Exposed: exposed,
141 Callbacks: mgr.collectorCallbacks,
142 WaitKey: func(cfg confgroup.Config) string {
143 return cfg.FullName()
144 },
145
146 Path: fmt.Sprintf(dyncfgCollectorPath, cfg.PluginName),
147 EnableFailCode: 200,
148 RemoveStockOnEnableFail: true,
149 JobCommands: []dyncfg.Command{
150 dyncfg.CommandSchema,
151 dyncfg.CommandGet,
152 dyncfg.CommandEnable,
153 dyncfg.CommandDisable,
154 dyncfg.CommandUpdate,
155 dyncfg.CommandRestart,
156 dyncfg.CommandTest,
157 dyncfg.CommandUserconfig,
158 },
159 })
160 mgr.vnodesCtl = vnodectl.New(vnodectl.Options{
161 Logger: mgr.Logger,
162 API: api,
163 Plugin: cfg.PluginName,
164 Initial: vnodesReg,
165 AffectedJobs: mgr.affectedVnodeJobs,
166 ApplyVnodeUpdate: mgr.applyVnodeUpdate,
167 })
168 mgr.secretsCtl = secretsctl.New(secretsctl.Options{
169 Logger: mgr.Logger,
170 API: mgr.dyncfgResponder,
171 Service: secretStoreSvc,
172 Plugin: mgr.pluginName,
173 Initial: mgr.initialSecretStores,
174 AffectedJobs: mgr.affectedJobs,
175 RestartableAffectedJobs: mgr.restartableAffectedJobs,
176 RestartDependentJobs: mgr.restartDependentJobs,
177 })
178
179 return mgr
180 }
181
182 // SetDyncfgResponder allows overriding the default responder (e.g., to silence output in CLI mode).
183 func (m *Manager) SetDyncfgResponder(responder *dyncfg.Responder) {
184 if responder != nil && m.dyncfgResponder != nil {
185 responder.SetTerminalFinalizer(m.dyncfgResponder.TerminalFinalizer())
186 }
187 dyncfg.BindResponder(&m.dyncfgResponder, m.collectorHandler, responder)
188 m.secretsCtl.SetAPI(responder)
189 m.vnodesCtl.SetAPI(responder)
190 m.funcCtl.SetAPI(responder)
191 }
192
193 type Manager struct {
194 *logger.Logger
195
196 // Static configuration and injected dependencies.
197 pluginName string
198 out io.Writer
199 runModePolicy policy.RunModePolicy
200 modules collectorapi.Registry
201 runJobNames []string
202 configDefaults confgroup.Registry
203 varLibDir string
204 fnReg FunctionRegistry
205 initialSecretStores []secretstore.Config
206
207 // Metrics-audit mode.
208 auditMode bool
209 auditAnalyzer metricsaudit.Analyzer
210 auditDataDir string
211
212 // Persistent caches and runtime state.
213 fileStatus *fileStatus
214 discoveredConfigs *discoveredConfigs
215 collectorSeen *dyncfg.SeenCache[confgroup.Config]
216 collectorExposed *dyncfg.ExposedCache[confgroup.Config]
217 secretStoreDeps *secretStoreDeps
218 retryingTasks *retryingTasks
219 runningJobs *runningJobs
220
221 // Controllers and handlers.
222 funcCtl *funcctl.Controller
223 collectorHandler *dyncfg.Handler[confgroup.Config]
224 collectorCallbacks *collectorCallbacks
225 secretsCtl *secretsctl.Controller
226 vnodesCtl *vnodectl.Controller
227
228 // Runtime loop state.
229 ctx context.Context
230 started chan struct{}
231 addCh chan confgroup.Config
232 rmCh chan confgroup.Config
233 dyncfgCh chan dyncfg.Function
234 cmdTestSem chan struct{}
235 cmdTestWG sync.WaitGroup
236
237 // Shared service seams.
238 dyncfgResponder *dyncfg.Responder
239
240 // RuntimeService is an optional runtime/internal metrics registration seam.
241 // When set, V2 jobs may register per-job runtime components.
242 runtimeService runtimecomp.Service
243 vnodeRegistry *vnoderegistry.Registry
244
245 secretResolver *secretresolver.Resolver
246 }
247
248 func (m *Manager) Run(ctx context.Context, in chan []*confgroup.Group) {
249 m.Info("instance is started")
250 defer func() { m.cleanup(); m.Info("instance is stopped") }()
251 m.ctx = ctx
252 m.funcCtl.Init(ctx)
253
254 vnodePrefix := m.vnodesCtl.Prefix()
255 m.fnReg.RegisterPrefix("config", vnodePrefix, dyncfg.WrapHandler(m.dyncfgConfig))
256 m.fnReg.RegisterPrefix("config", m.dyncfgSecretStorePrefixValue(), dyncfg.WrapHandler(m.dyncfgConfig))
257
258 m.vnodesCtl.CreateTemplates()
259 m.vnodesCtl.PublishExisting(dyncfg.StatusRunning)
260
261 m.secretsCtl.CreateTemplates()
262 m.secretsCtl.PublishExisting()
263
264 m.fnReg.RegisterPrefix("config", m.dyncfgCollectorPrefixValue(), dyncfg.WrapHandler(m.dyncfgConfig))
265 for name := range m.modules {
266 m.dyncfgCollectorModuleCreate(name)
267 }
268 m.funcCtl.RegisterModules(m.modules)
269
270 m.loadFileStatus()
271
272 var wg sync.WaitGroup
273
274 wg.Go(func() { m.runFileStatusPersistence() })
275
276 wg.Go(func() { m.runProcessConfGroups(in) })
277
278 wg.Go(func() { m.run() })
279
280 wg.Go(func() { m.runNotifyRunningJobs() })
281
282 close(m.started)
283
284 wg.Wait()
285 <-m.ctx.Done()
286 }
287
288 // WaitStarted blocks until Run has completed initialization or the context is canceled.
289 func (m *Manager) WaitStarted(ctx context.Context) bool {
290 select {
291 case <-m.started:
292 return true
293 case <-ctx.Done():
294 return false
295 }
296 }
297
298 // GetJobNames returns the currently running job names for a module.
299 func (m *Manager) GetJobNames(moduleName string) []string {
300 return m.funcCtl.GetJobNames(moduleName)
301 }
302
303 // ExecuteFunction executes a function handler directly (function name must be module:method).
304 func (m *Manager) ExecuteFunction(functionName string, fn functions.Function) {
305 m.funcCtl.ExecuteFunction(functionName, fn)
306 }
307
308 func (m *Manager) runProcessConfGroups(in chan []*confgroup.Group) {
309 for {
310 select {
311 case <-m.ctx.Done():
312 return
313 case groups, ok := <-in:
314 if !ok {
315 return
316 }
317 for _, gr := range groups {
318 a, r := m.discoveredConfigs.add(gr)
319 m.Debugf("received configs: %d/+%d/-%d ('%s')", len(gr.Configs), len(a), len(r), gr.Source)
320 if len(m.runJobNames) > 0 {
321 a = slices.DeleteFunc(a, func(config confgroup.Config) bool {
322 return !slices.ContainsFunc(m.runJobNames, func(name string) bool { return config.Name() == name })
323 })
324 }
325 sendConfigs(m.ctx, m.rmCh, r...)
326 sendConfigs(m.ctx, m.addCh, a...)
327 }
328 }
329 }
330 }
331
332 func (m *Manager) run() {
333 for {
334 if m.collectorHandler.WaitingForDecision() {
335 step, ok := m.collectorHandler.NextWaitDecisionStep(m.ctx, m.dyncfgCh)
336 if !ok {
337 return
338 }
339 if step.HasCommand {
340 m.dyncfgSeqExec(step.Command)
341 continue
342 }
343 } else {
344 select {
345 case <-m.ctx.Done():
346 return
347 case cfg := <-m.addCh:
348 m.addConfig(cfg)
349 case cfg := <-m.rmCh:
350 m.removeConfig(cfg)
351 case fn := <-m.dyncfgCh:
352 m.dyncfgSeqExec(fn)
353 }
354 }
355 }
356 }
357
358 func (m *Manager) addConfig(cfg confgroup.Config) {
359 if _, ok := m.modules.Lookup(cfg.Module()); !ok {
360 return
361 }
362
363 m.retryingTasks.remove(cfg)
364
365 m.collectorHandler.RememberDiscoveredConfig(cfg)
366
367 entry, ok := m.collectorExposed.LookupByKey(cfg.ExposedKey())
368 if !ok {
369 entry = m.collectorHandler.AddDiscoveredConfig(cfg, dyncfg.StatusAccepted)
370 } else {
371 sp, ep := cfg.SourceTypePriority(), entry.Cfg.SourceTypePriority()
372 if ep > sp || (ep == sp && entry.Status == dyncfg.StatusRunning) {
373 return
374 }
375 if entry.Status == dyncfg.StatusRunning {
376 m.stopRunningJob(entry.Cfg.FullName())
377 m.fileStatus.remove(entry.Cfg)
378 }
379 entry = m.collectorHandler.AddDiscoveredConfig(cfg, dyncfg.StatusAccepted) // replace existing exposed
380 }
381
382 m.syncSecretStoreDepsForConfig(entry.Cfg)
383 m.collectorHandler.NotifyJobCreate(entry.Cfg, entry.Status)
384
385 if m.runModePolicy.AutoEnableDiscovered {
386 m.collectorHandler.CmdEnable(dyncfg.NewFunction(functions.Function{Args: []string{m.dyncfgJobID(entry.Cfg), "enable"}}))
387 } else {
388 m.collectorHandler.WaitForDecision(entry.Cfg)
389 }
390 }
391
392 func (m *Manager) removeConfig(cfg confgroup.Config) {
393 m.retryingTasks.remove(cfg)
394
395 entry, ok := m.collectorHandler.RemoveDiscoveredConfig(cfg)
396 if !ok {
397 return
398 }
399 // Stop first so running-state cleanup is applied against existing dependency state.
400 m.stopRunningJob(cfg.FullName())
401 m.secretStoreDeps.RemoveActiveJob(entry.Cfg.FullName())
402 m.fileStatus.remove(cfg)
403
404 if !isStock(cfg) || entry.Status == dyncfg.StatusRunning {
405 m.collectorHandler.NotifyJobRemove(cfg)
406 }
407 }
408
409 func (m *Manager) runNotifyRunningJobs() {
410 tk := ticker.New(time.Second)
411 defer tk.Stop()
412
413 for {
414 select {
415 case <-m.ctx.Done():
416 return
417 case clock := <-tk.C:
418 for _, job := range m.runningJobs.snapshot() {
419 job.Tick(clock)
420 }
421 }
422 }
423 }
424
425 func (m *Manager) startRunningJob(job runtimeJob) {
426 m.stopRunningJob(job.FullName())
427
428 go job.Start()
429
430 m.runningJobs.lock()
431 m.runningJobs.add(job.FullName(), job)
432 m.runningJobs.unlock()
433 m.secretStoreDeps.setRunning(job.FullName(), true)
434
435 // Known behavior: Start runs asynchronously, so function handlers may be
436 // published before the job flips its running flag. Immediate function calls
437 // can transiently return 503; this is accepted for now to avoid adding a
438 // broader runtime readiness contract.
439 m.funcCtl.OnJobStart(job)
440 }
441
442 func (m *Manager) stopRunningJob(name string) {
443 m.runningJobs.lock()
444 job, ok := m.runningJobs.lookup(name)
445 if ok {
446 m.runningJobs.remove(name)
447 }
448 m.runningJobs.unlock()
449 if ok {
450 m.secretStoreDeps.setRunning(name, false)
451 m.funcCtl.OnJobStop(job)
452 job.Stop()
453 }
454 }
455
456 func (m *Manager) cleanup() {
457 m.fnReg.UnregisterPrefix("config", m.dyncfgCollectorPrefixValue())
458 m.fnReg.UnregisterPrefix("config", m.dyncfgSecretStorePrefixValue())
459 m.fnReg.UnregisterPrefix("config", m.dyncfgVnodePrefixValue())
460 m.funcCtl.Cleanup()
461
462 for _, job := range m.runningJobs.snapshot() {
463 m.stopRunningJob(job.FullName())
464 }
465
466 m.waitCmdTestWorkers()
467 }
468
469 func (m *Manager) waitCmdTestWorkers() {
470 done := make(chan struct{})
471 go func() {
472 m.cmdTestWG.Wait()
473 close(done)
474 }()
475
476 select {
477 case <-done:
478 case <-time.After(cmdTestWorkerDrainWait):
479 m.Warningf("dyncfg: timeout waiting %s for command test workers to drain", cmdTestWorkerDrainWait)
480 }
481 }
482
483 func (m *Manager) createCollectorJob(cfg confgroup.Config) (runtimeJob, error) {
484 return newJobFactory(m).create(cfg)
485 }
486
487 func (m *Manager) validateCollectorJob(cfg confgroup.Config) error {
488 return newJobFactory(m).validate(cfg)
489 }
490
491 func (m *Manager) baseContext() context.Context {
492 if m.ctx != nil {
493 return m.ctx
494 }
495 return context.Background()
496 }
497
498 func runRetryTask(ctx context.Context, out chan<- confgroup.Config, cfg confgroup.Config) {
499 t := time.NewTimer(time.Second * time.Duration(cfg.AutoDetectionRetry()))
500 defer t.Stop()
501
502 select {
503 case <-ctx.Done():
504 case <-t.C:
505 sendConfigs(ctx, out, cfg)
506 }
507 }
508
509 func sendConfigs(ctx context.Context, out chan<- confgroup.Config, cfgs ...confgroup.Config) {
510 for _, cfg := range cfgs {
511 select {
512 case <-ctx.Done():
513 return
514 case out <- cfg:
515 }
516 }
517 }
518
519 func isStock(cfg confgroup.Config) bool {
520 return cfg.SourceType() == confgroup.TypeStock
521 }
522
523 func isDyncfg(cfg confgroup.Config) bool {
524 return cfg.SourceType() == confgroup.TypeDyncfg
525 }