master
go 220 lines 6.51 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 "os"
10 "path/filepath"
11
12 "github.com/netdata/netdata/go/plugins/logger"
13 "github.com/netdata/netdata/go/plugins/plugin/agent/internal/naming"
14 "github.com/netdata/netdata/go/plugins/plugin/agent/secrets/resolver"
15 "github.com/netdata/netdata/go/plugins/plugin/agent/secrets/secretstore"
16 "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi"
17 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
18 "github.com/netdata/netdata/go/plugins/plugin/framework/jobruntime"
19 "github.com/netdata/netdata/go/plugins/plugin/framework/metricsaudit"
20 "github.com/netdata/netdata/go/plugins/plugin/framework/runtimecomp"
21 "github.com/netdata/netdata/go/plugins/plugin/framework/vnoderegistry"
22 "github.com/netdata/netdata/go/plugins/plugin/framework/vnodes"
23 )
24
25 func jobLogSource(cfg confgroup.Config) string {
26 sourceType := cfg.SourceType()
27 provider := cfg.Provider()
28 if sourceType != "" && sourceType == provider {
29 return sourceType
30 }
31 return fmt.Sprintf("%s/%s", sourceType, provider)
32 }
33
34 // jobFactory builds runtime jobs from configs without mutating manager-owned runtime maps.
35 type jobFactory struct {
36 logger *logger.Logger
37
38 pluginName string
39 modules collectorapi.Registry
40 vnodeLookup func(string) (*vnodes.VirtualNode, bool)
41 out io.Writer
42
43 validationOnly bool
44
45 auditMode bool
46 auditAnalyzer metricsaudit.Analyzer
47 auditDataDir string
48
49 runtimeService runtimecomp.Service
50 vnodeRegistry *vnoderegistry.Registry
51
52 secretResolver *secretresolver.Resolver
53 secretStoreSvc secretstore.Service
54 ctx context.Context
55 }
56
57 func newJobFactory(m *Manager) *jobFactory {
58 return &jobFactory{
59 logger: m.Logger,
60
61 pluginName: m.pluginName,
62 modules: m.modules,
63 vnodeLookup: m.vnodesCtl.Lookup,
64 out: m.out,
65
66 auditMode: m.auditMode,
67 auditAnalyzer: m.auditAnalyzer,
68 auditDataDir: m.auditDataDir,
69
70 runtimeService: m.runtimeService,
71 vnodeRegistry: m.vnodeRegistry,
72 secretResolver: m.secretResolver,
73 secretStoreSvc: m.secretsCtl.Service(),
74 ctx: m.baseContext(),
75 }
76 }
77
78 func (f *jobFactory) validate(cfg confgroup.Config) error {
79 clone := *f
80 clone.validationOnly = true
81 _, err := clone.create(cfg)
82 return err
83 }
84
85 func (f *jobFactory) create(cfg confgroup.Config) (runtimeJob, error) {
86 creator, ok := f.modules[cfg.Module()]
87 if !ok {
88 return nil, fmt.Errorf("can not find %s module", cfg.Module())
89 }
90
91 functionOnly := creator.FunctionOnly || cfg.FunctionOnly()
92 if cfg.FunctionOnly() && creator.Methods == nil && creator.JobMethods == nil {
93 return nil, fmt.Errorf("function_only is set but %s module has no methods defined", cfg.Module())
94 }
95
96 var vnode *vnodes.VirtualNode
97 if cfg.Vnode() != "" {
98 if f.vnodeLookup == nil {
99 return nil, fmt.Errorf("vnode '%s' is not found", cfg.Vnode())
100 }
101 n, ok := f.vnodeLookup(cfg.Vnode())
102 if !ok || n == nil {
103 return nil, fmt.Errorf("vnode '%s' is not found", cfg.Vnode())
104 }
105 vnode = n
106 }
107
108 f.logger.Debugf("creating %s[%s] job, config: %v", cfg.Module(), cfg.Name(), cfg)
109
110 if creator.CreateV2 != nil {
111 return f.createV2(cfg, creator, functionOnly, vnode)
112 }
113 return f.createV1(cfg, creator, functionOnly, vnode)
114 }
115
116 func (f *jobFactory) logApplyConfigError(cfg confgroup.Config, err error) {
117 if f.validationOnly {
118 return
119 }
120 f.logger.Errorf("failed to apply config for %s[%s] job: %v", cfg.Module(), cfg.Name(), err)
121 }
122
123 func (f *jobFactory) createV2(cfg confgroup.Config, creator collectorapi.Creator, functionOnly bool, vnode *vnodes.VirtualNode) (runtimeJob, error) {
124 mod := creator.CreateV2()
125 if mod == nil {
126 return nil, fmt.Errorf("module %s CreateV2 returned nil", cfg.Module())
127 }
128 storeSnapshot := f.secretStoreSvc.Capture()
129 resolveCtx := collectorSecretResolveContext(f.ctx, f.logger, cfg)
130 if err := applyConfig(resolveCtx, cfg, mod, f.secretResolver, f.secretStoreSvc, storeSnapshot); err != nil {
131 f.logApplyConfigError(cfg, err)
132 return nil, err
133 }
134
135 jobCfg := jobruntime.JobV2Config{
136 PluginName: f.pluginName,
137 Name: cfg.Name(),
138 ModuleName: cfg.Module(),
139 FullName: cfg.FullName(),
140 Source: jobLogSource(cfg),
141 UpdateEvery: cfg.UpdateEvery(),
142 AutoDetectEvery: cfg.AutoDetectionRetry(),
143 IsStock: cfg.SourceType() == "stock",
144 Labels: makeLabels(cfg),
145 Out: f.out,
146 Module: mod,
147 FunctionOnly: functionOnly,
148 RuntimeService: f.runtimeService,
149 VnodeRegistry: f.vnodeRegistry,
150 }
151 if vnode != nil {
152 jobCfg.Vnode = *vnode.Copy()
153 }
154 return jobruntime.NewJobV2(jobCfg), nil
155 }
156
157 func (f *jobFactory) createV1(cfg confgroup.Config, creator collectorapi.Creator, functionOnly bool, vnode *vnodes.VirtualNode) (runtimeJob, error) {
158 if creator.Create == nil {
159 return nil, fmt.Errorf("module %s has no compatible creator", cfg.Module())
160 }
161
162 jobCaptureDir, err := f.createV1CaptureDir(cfg)
163 if err != nil {
164 return nil, err
165 }
166
167 mod := creator.Create()
168 storeSnapshot := f.secretStoreSvc.Capture()
169 resolveCtx := collectorSecretResolveContext(f.ctx, f.logger, cfg)
170 if err := applyConfig(resolveCtx, cfg, mod, f.secretResolver, f.secretStoreSvc, storeSnapshot); err != nil {
171 f.logApplyConfigError(cfg, err)
172 return nil, err
173 }
174
175 if f.auditAnalyzer != nil && jobCaptureDir != "" {
176 f.auditAnalyzer.RegisterJob(cfg.Name(), cfg.Module(), jobCaptureDir)
177 }
178 if jobCaptureDir != "" {
179 if captureAware, ok := mod.(metricsaudit.Capturable); ok {
180 captureAware.EnableCaptureArtifacts(jobCaptureDir)
181 }
182 }
183
184 jobCfg := jobruntime.JobConfig{
185 PluginName: f.pluginName,
186 Name: cfg.Name(),
187 ModuleName: cfg.Module(),
188 FullName: cfg.FullName(),
189 Source: jobLogSource(cfg),
190 UpdateEvery: cfg.UpdateEvery(),
191 AutoDetectEvery: cfg.AutoDetectionRetry(),
192 Priority: cfg.Priority(),
193 Labels: makeLabels(cfg),
194 IsStock: cfg.SourceType() == "stock",
195 Module: mod,
196 Out: f.out,
197 AuditMode: f.auditMode,
198 AuditAnalyzer: f.auditAnalyzer,
199 FunctionOnly: functionOnly,
200 }
201 if vnode != nil {
202 jobCfg.Vnode = *vnode.Copy()
203 }
204
205 return jobruntime.NewJob(jobCfg), nil
206 }
207
208 func (f *jobFactory) createV1CaptureDir(cfg confgroup.Config) (string, error) {
209 if f.validationOnly {
210 return "", nil
211 }
212 if f.auditDataDir == "" {
213 return "", nil
214 }
215 jobCaptureDir := filepath.Join(f.auditDataDir, naming.Sanitize(cfg.Module()), naming.Sanitize(cfg.Name()))
216 if err := os.MkdirAll(jobCaptureDir, 0o755); err != nil {
217 return "", fmt.Errorf("creating audit directory: %w", err)
218 }
219 return jobCaptureDir, nil
220 }